dht.go 16.9 KB
Newer Older
Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
1 2
package dht

3
import (
4
	"bytes"
5 6 7
	"errors"
	"sync"
	"time"
8

9
	peer "github.com/jbenet/go-ipfs/peer"
Jeromy's avatar
Jeromy committed
10
	kb "github.com/jbenet/go-ipfs/routing/kbucket"
11 12
	swarm "github.com/jbenet/go-ipfs/swarm"
	u "github.com/jbenet/go-ipfs/util"
13 14

	ma "github.com/jbenet/go-multiaddr"
Jeromy's avatar
Jeromy committed
15 16 17

	ds "github.com/jbenet/datastore.go"

18
	"code.google.com/p/goprotobuf/proto"
19 20
)

Juan Batiz-Benet's avatar
Juan Batiz-Benet committed
21 22 23 24 25
// TODO. SEE https://github.com/jbenet/node-ipfs/blob/master/submodules/ipfs-dht/index.js

// IpfsDHT is an implementation of Kademlia with Coral and S/Kademlia modifications.
// It is used to implement the base IpfsRouting module.
type IpfsDHT struct {
26 27
	// Array of routing tables for differently distanced nodes
	// NOTE: (currently, only a single table is used)
28
	routes []*kb.RoutingTable
29

30
	network swarm.Network
31

Jeromy's avatar
Jeromy committed
32 33 34 35 36
	// Local peer (yourself)
	self *peer.Peer

	// Local data
	datastore ds.Datastore
37

38
	// Map keys to peers that can provide their value
39
	providers    map[u.Key][]*providerInfo
40
	providerLock sync.RWMutex
41

42 43
	// Signal to shutdown dht
	shutdown chan struct{}
44 45 46

	// When this peer started up
	birth time.Time
47 48 49

	//lock to make diagnostics work better
	diaglock sync.Mutex
50

51 52
	// listener is a server to register to listen for responses to messages
	listener *MesListener
53 54
}

Jeromy's avatar
Jeromy committed
55
// NewDHT creates a new DHT object with the given peer as the 'local' host
Jeromy's avatar
Jeromy committed
56
func NewDHT(p *peer.Peer, net swarm.Network) *IpfsDHT {
57
	dht := new(IpfsDHT)
Jeromy's avatar
Jeromy committed
58
	dht.network = net
59 60
	dht.datastore = ds.NewMapDatastore()
	dht.self = p
61
	dht.providers = make(map[u.Key][]*providerInfo)
Jeromy's avatar
Jeromy committed
62
	dht.shutdown = make(chan struct{})
63 64 65 66 67 68

	dht.routes = make([]*kb.RoutingTable, 3)
	dht.routes[0] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*30)
	dht.routes[1] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Millisecond*100)
	dht.routes[2] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID), time.Hour)

69
	dht.listener = NewMesListener()
70
	dht.birth = time.Now()
Jeromy's avatar
Jeromy committed
71
	return dht
72 73
}

74
// Start up background goroutines needed by the DHT
75 76 77 78
func (dht *IpfsDHT) Start() {
	go dht.handleMessages()
}

79
// Connect to a new peer at the given address, ping and add to the routing table
80
func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
81
	maddrstr, _ := addr.String()
82
	u.DOut("Connect to new peer: %s", maddrstr)
83
	npeer, err := dht.network.ConnectNew(addr)
84
	if err != nil {
85
		return nil, err
86 87
	}

Jeromy's avatar
Jeromy committed
88 89
	// Ping new peer to register in their routing table
	// NOTE: this should be done better...
90
	err = dht.Ping(npeer, time.Second*2)
Jeromy's avatar
Jeromy committed
91
	if err != nil {
Jeromy's avatar
Jeromy committed
92
		return nil, errors.New("failed to ping newly connected peer")
Jeromy's avatar
Jeromy committed
93 94
	}

95 96
	dht.Update(npeer)

97
	return npeer, nil
Jeromy's avatar
Jeromy committed
98 99
}

100 101
// Read in all messages from swarm and handle them appropriately
// NOTE: this function is just a quick sketch
102
func (dht *IpfsDHT) handleMessages() {
103
	u.DOut("Begin message handling routine")
104 105

	checkTimeouts := time.NewTicker(time.Minute * 5)
106
	ch := dht.network.GetChan()
107 108
	for {
		select {
109
		case mes, ok := <-ch.Incoming:
110 111 112 113
			if !ok {
				u.DOut("handleMessages closing, bad recv on incoming")
				return
			}
114
			pmes := new(PBDHTMessage)
115 116 117 118 119 120
			err := proto.Unmarshal(mes.Data, pmes)
			if err != nil {
				u.PErr("Failed to decode protobuf message: %s", err)
				continue
			}

121
			dht.Update(mes.Peer)
122

123
			// Note: not sure if this is the correct place for this
124
			if pmes.GetResponse() {
125
				dht.listener.Respond(pmes.GetId(), mes)
126
				continue
127
			}
128 129
			//

Jeromy's avatar
Jeromy committed
130 131
			u.DOut("[peer: %s]\nGot message type: '%s' [id = %x, from = %s]",
				dht.self.ID.Pretty(),
132
				PBDHTMessage_MessageType_name[int32(pmes.GetType())],
133
				pmes.GetId(), mes.Peer.ID.Pretty())
134
			switch pmes.GetType() {
135
			case PBDHTMessage_GET_VALUE:
136
				dht.handleGetValue(mes.Peer, pmes)
137
			case PBDHTMessage_PUT_VALUE:
Jeromy's avatar
Jeromy committed
138
				dht.handlePutValue(mes.Peer, pmes)
139
			case PBDHTMessage_FIND_NODE:
Jeromy's avatar
Jeromy committed
140
				dht.handleFindPeer(mes.Peer, pmes)
141
			case PBDHTMessage_ADD_PROVIDER:
142
				dht.handleAddProvider(mes.Peer, pmes)
143
			case PBDHTMessage_GET_PROVIDERS:
144
				dht.handleGetProviders(mes.Peer, pmes)
145
			case PBDHTMessage_PING:
146
				dht.handlePing(mes.Peer, pmes)
147
			case PBDHTMessage_DIAGNOSTIC:
148
				dht.handleDiagnostic(mes.Peer, pmes)
149 150
			}

151 152
		case err := <-ch.Errors:
			u.PErr("dht err: %s", err)
153
		case <-dht.shutdown:
154
			checkTimeouts.Stop()
155
			return
156
		case <-checkTimeouts.C:
157 158 159 160 161 162 163 164 165 166 167 168 169
			// Time to collect some garbage!
			dht.cleanExpiredProviders()
		}
	}
}

func (dht *IpfsDHT) cleanExpiredProviders() {
	dht.providerLock.Lock()
	for k, parr := range dht.providers {
		var cleaned []*providerInfo
		for _, v := range parr {
			if time.Since(v.Creation) < time.Hour {
				cleaned = append(cleaned, v)
170
			}
171
		}
172
		dht.providers[k] = cleaned
173
	}
174 175 176
	dht.providerLock.Unlock()
}

Jeromy's avatar
Jeromy committed
177
func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) error {
178 179 180
	pmes := DHTMessage{
		Type:  PBDHTMessage_PUT_VALUE,
		Key:   key,
181
		Value: value,
182
		Id:    GenerateMessageID(),
183 184 185
	}

	mes := swarm.NewMessage(p, pmes.ToProtobuf())
186
	dht.network.Send(mes)
187 188 189
	return nil
}

190
func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
Jeromy's avatar
Jeromy committed
191
	u.DOut("handleGetValue for key: %s", pmes.GetKey())
Jeromy's avatar
Jeromy committed
192
	dskey := ds.NewKey(pmes.GetKey())
Jeromy's avatar
Jeromy committed
193 194 195 196 197 198
	resp := &DHTMessage{
		Response: true,
		Id:       pmes.GetId(),
		Key:      pmes.GetKey(),
	}
	iVal, err := dht.datastore.Get(dskey)
Jeromy's avatar
Jeromy committed
199
	if err == nil {
Jeromy's avatar
Jeromy committed
200
		u.DOut("handleGetValue success!")
Jeromy's avatar
Jeromy committed
201 202
		resp.Success = true
		resp.Value = iVal.([]byte)
Jeromy's avatar
Jeromy committed
203
	} else if err == ds.ErrNotFound {
Jeromy's avatar
Jeromy committed
204 205 206
		// Check if we know any providers for the requested value
		provs, ok := dht.providers[u.Key(pmes.GetKey())]
		if ok && len(provs) > 0 {
Jeromy's avatar
Jeromy committed
207
			u.DOut("handleGetValue returning %d provider[s]", len(provs))
Jeromy's avatar
Jeromy committed
208 209 210 211 212 213
			for _, prov := range provs {
				resp.Peers = append(resp.Peers, prov.Value)
			}
			resp.Success = true
		} else {
			// No providers?
214 215
			// Find closest peer on given cluster to desired key and reply with that info

216 217 218 219 220 221 222
			level := 0
			if len(pmes.GetValue()) < 1 {
				// TODO: maybe return an error? Defaulting isnt a good idea IMO
				u.PErr("handleGetValue: no routing level specified, assuming 0")
			} else {
				level = int(pmes.GetValue()[0]) // Using value field to specify cluster level
			}
Jeromy's avatar
Jeromy committed
223
			u.DOut("handleGetValue searching level %d clusters", level)
224 225 226

			closer := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))

Jeromy's avatar
Jeromy committed
227 228 229 230 231
			if closer.ID.Equal(dht.self.ID) {
				u.DOut("Attempted to return self! this shouldnt happen...")
				resp.Peers = nil
				goto out
			}
232 233 234
			// If this peer is closer than the one from the table, return nil
			if kb.Closer(dht.self.ID, closer.ID, u.Key(pmes.GetKey())) {
				resp.Peers = nil
Jeromy's avatar
Jeromy committed
235
				u.DOut("handleGetValue could not find a closer node than myself.")
236
			} else {
Jeromy's avatar
Jeromy committed
237
				u.DOut("handleGetValue returning a closer peer: '%s'", closer.ID.Pretty())
238 239
				resp.Peers = []*peer.Peer{closer}
			}
240
		}
Jeromy's avatar
Jeromy committed
241
	} else {
242
		//temp: what other errors can a datastore return?
Jeromy's avatar
Jeromy committed
243
		panic(err)
244
	}
245

Jeromy's avatar
Jeromy committed
246
out:
247
	mes := swarm.NewMessage(p, resp.ToProtobuf())
248
	dht.network.Send(mes)
249 250
}

Jeromy's avatar
Jeromy committed
251
// Store a value in this peer local storage
252
func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *PBDHTMessage) {
Jeromy's avatar
Jeromy committed
253 254 255 256 257 258
	dskey := ds.NewKey(pmes.GetKey())
	err := dht.datastore.Put(dskey, pmes.GetValue())
	if err != nil {
		// For now, just panic, handle this better later maybe
		panic(err)
	}
259 260
}

261 262 263
func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *PBDHTMessage) {
	resp := DHTMessage{
		Type:     pmes.GetType(),
264
		Response: true,
265
		Id:       pmes.GetId(),
266
	}
267

268
	dht.network.Send(swarm.NewMessage(p, resp.ToProtobuf()))
269 270
}

271
func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *PBDHTMessage) {
272 273 274 275 276 277 278 279 280 281 282 283
	resp := DHTMessage{
		Type:     pmes.GetType(),
		Id:       pmes.GetId(),
		Response: true,
	}
	defer func() {
		mes := swarm.NewMessage(p, resp.ToProtobuf())
		dht.network.Send(mes)
	}()
	level := pmes.GetValue()[0]
	u.DOut("handleFindPeer: searching for '%s'", peer.ID(pmes.GetKey()).Pretty())
	closest := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
Jeromy's avatar
Jeromy committed
284
	if closest == nil {
285
		u.PErr("handleFindPeer: could not find anything.")
286
		return
Jeromy's avatar
Jeromy committed
287 288 289
	}

	if len(closest.Addresses) == 0 {
290
		u.PErr("handleFindPeer: no addresses for connected peer...")
291
		return
Jeromy's avatar
Jeromy committed
292 293
	}

294 295 296
	// If the found peer further away than this peer...
	if kb.Closer(dht.self.ID, closest.ID, u.Key(pmes.GetKey())) {
		return
Jeromy's avatar
Jeromy committed
297 298
	}

299 300 301 302 303 304
	u.DOut("handleFindPeer: sending back '%s'", closest.ID.Pretty())
	resp.Peers = []*peer.Peer{closest}
	resp.Success = true
}

func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *PBDHTMessage) {
305
	resp := DHTMessage{
306 307
		Type:     PBDHTMessage_GET_PROVIDERS,
		Key:      pmes.GetKey(),
308
		Id:       pmes.GetId(),
309
		Response: true,
Jeromy's avatar
Jeromy committed
310 311
	}

312
	dht.providerLock.RLock()
313
	providers := dht.providers[u.Key(pmes.GetKey())]
314
	dht.providerLock.RUnlock()
315
	if providers == nil || len(providers) == 0 {
Jeromy's avatar
Jeromy committed
316 317 318 319 320 321 322 323 324 325 326
		level := 0
		if len(pmes.GetValue()) > 0 {
			level = int(pmes.GetValue()[0])
		}

		closer := dht.routes[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
		if kb.Closer(dht.self.ID, closer.ID, u.Key(pmes.GetKey())) {
			resp.Peers = nil
		} else {
			resp.Peers = []*peer.Peer{closer}
		}
327 328 329
	} else {
		for _, prov := range providers {
			resp.Peers = append(resp.Peers, prov.Value)
330
		}
331
		resp.Success = true
332 333 334
	}

	mes := swarm.NewMessage(p, resp.ToProtobuf())
335
	dht.network.Send(mes)
336 337
}

338 339
type providerInfo struct {
	Creation time.Time
340
	Value    *peer.Peer
341 342
}

343
func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *PBDHTMessage) {
344 345
	//TODO: need to implement TTLs on providers
	key := u.Key(pmes.GetKey())
346
	dht.addProviderEntry(key, p)
347 348
}

Jeromy's avatar
Jeromy committed
349
// Stop all communications from this peer and shut down
350 351 352 353
func (dht *IpfsDHT) Halt() {
	dht.shutdown <- struct{}{}
	dht.network.Close()
}
354

355 356 357 358
func (dht *IpfsDHT) addProviderEntry(key u.Key, p *peer.Peer) {
	u.DOut("Adding %s as provider for '%s'", p.Key().Pretty(), key)
	dht.providerLock.Lock()
	provs := dht.providers[key]
359
	dht.providers[key] = append(provs, &providerInfo{time.Now(), p})
360 361
	dht.providerLock.Unlock()
}
362

Jeromy's avatar
Jeromy committed
363
// NOTE: not yet finished, low priority
364
func (dht *IpfsDHT) handleDiagnostic(p *peer.Peer, pmes *PBDHTMessage) {
365
	seq := dht.routes[0].NearestPeers(kb.ConvertPeerID(dht.self.ID), 10)
366
	listenChan := dht.listener.Listen(pmes.GetId(), len(seq), time.Second*30)
367

368
	for _, ps := range seq {
369
		mes := swarm.NewMessage(ps, pmes)
370
		dht.network.Send(mes)
371 372 373
	}

	buf := new(bytes.Buffer)
374 375 376
	di := dht.getDiagInfo()
	buf.Write(di.Marshal())

377 378 379 380 381 382 383 384
	// NOTE: this shouldnt be a hardcoded value
	after := time.After(time.Second * 20)
	count := len(seq)
	for count > 0 {
		select {
		case <-after:
			//Timeout, return what we have
			goto out
Jeromy's avatar
Jeromy committed
385
		case req_resp := <-listenChan:
386
			pmes_out := new(PBDHTMessage)
387 388 389 390 391
			err := proto.Unmarshal(req_resp.Data, pmes_out)
			if err != nil {
				// It broke? eh, whatever, keep going
				continue
			}
392 393 394 395 396 397
			buf.Write(req_resp.Data)
			count--
		}
	}

out:
398 399 400 401
	resp := DHTMessage{
		Type:     PBDHTMessage_DIAGNOSTIC,
		Id:       pmes.GetId(),
		Value:    buf.Bytes(),
402 403 404 405
		Response: true,
	}

	mes := swarm.NewMessage(p, resp.ToProtobuf())
406
	dht.network.Send(mes)
407
}
408

409 410 411
func (dht *IpfsDHT) getValueOrPeers(p *peer.Peer, key u.Key, timeout time.Duration, level int) ([]byte, []*peer.Peer, error) {
	pmes, err := dht.getValueSingle(p, key, timeout, level)
	if err != nil {
412
		return nil, nil, err
413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429
	}

	if pmes.GetSuccess() {
		if pmes.Value == nil { // We were given provider[s]
			val, err := dht.getFromPeerList(key, timeout, pmes.GetPeers(), level)
			if err != nil {
				return nil, nil, err
			}
			return val, nil, nil
		}

		// Success! We were given the value
		return pmes.GetValue(), nil, nil
	} else {
		// We were given a closer node
		var peers []*peer.Peer
		for _, pb := range pmes.GetPeers() {
430 431 432
			if peer.ID(pb.GetId()).Equal(dht.self.ID) {
				continue
			}
433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450
			addr, err := ma.NewMultiaddr(pb.GetAddr())
			if err != nil {
				u.PErr(err.Error())
				continue
			}

			np, err := dht.network.GetConnection(peer.ID(pb.GetId()), addr)
			if err != nil {
				u.PErr(err.Error())
				continue
			}

			peers = append(peers, np)
		}
		return nil, peers, nil
	}
}

451 452
// getValueSingle simply performs the get value RPC with the given parameters
func (dht *IpfsDHT) getValueSingle(p *peer.Peer, key u.Key, timeout time.Duration, level int) (*PBDHTMessage, error) {
Jeromy's avatar
Jeromy committed
453
	pmes := DHTMessage{
454 455 456 457
		Type:  PBDHTMessage_GET_VALUE,
		Key:   string(key),
		Value: []byte{byte(level)},
		Id:    GenerateMessageID(),
Jeromy's avatar
Jeromy committed
458
	}
459
	response_chan := dht.listener.Listen(pmes.Id, 1, time.Minute)
Jeromy's avatar
Jeromy committed
460 461

	mes := swarm.NewMessage(p, pmes.ToProtobuf())
462
	t := time.Now()
463
	dht.network.Send(mes)
Jeromy's avatar
Jeromy committed
464 465 466 467 468

	// Wait for either the response or a timeout
	timeup := time.After(timeout)
	select {
	case <-timeup:
469
		dht.listener.Unlisten(pmes.Id)
Jeromy's avatar
Jeromy committed
470 471 472 473 474 475
		return nil, u.ErrTimeout
	case resp, ok := <-response_chan:
		if !ok {
			u.PErr("response channel closed before timeout, please investigate.")
			return nil, u.ErrTimeout
		}
476 477
		roundtrip := time.Since(t)
		resp.Peer.SetLatency(roundtrip)
Jeromy's avatar
Jeromy committed
478 479 480 481 482
		pmes_out := new(PBDHTMessage)
		err := proto.Unmarshal(resp.Data, pmes_out)
		if err != nil {
			return nil, err
		}
483
		return pmes_out, nil
Jeromy's avatar
Jeromy committed
484 485 486
	}
}

487 488 489 490 491 492 493 494 495 496
// TODO: Im not certain on this implementation, we get a list of peers/providers
// from someone what do we do with it? Connect to each of them? randomly pick
// one to get the value from? Or just connect to one at a time until we get a
// successful connection and request the value from it?
func (dht *IpfsDHT) getFromPeerList(key u.Key, timeout time.Duration,
	peerlist []*PBDHTMessage_PBPeer, level int) ([]byte, error) {
	for _, pinfo := range peerlist {
		p, _ := dht.Find(peer.ID(pinfo.GetId()))
		if p == nil {
			maddr, err := ma.NewMultiaddr(pinfo.GetAddr())
Jeromy's avatar
Jeromy committed
497 498 499 500
			if err != nil {
				u.PErr("getValue error: %s", err)
				continue
			}
501

502
			p, err = dht.network.GetConnection(peer.ID(pinfo.GetId()), maddr)
Jeromy's avatar
Jeromy committed
503 504 505 506 507
			if err != nil {
				u.PErr("getValue error: %s", err)
				continue
			}
		}
508
		pmes, err := dht.getValueSingle(p, key, timeout, level)
Jeromy's avatar
Jeromy committed
509
		if err != nil {
510
			u.DErr("getFromPeers error: %s", err)
Jeromy's avatar
Jeromy committed
511 512
			continue
		}
513
		dht.addProviderEntry(key, p)
Jeromy's avatar
Jeromy committed
514

515 516 517 518
		// Make sure it was a successful get
		if pmes.GetSuccess() && pmes.Value != nil {
			return pmes.GetValue(), nil
		}
Jeromy's avatar
Jeromy committed
519 520 521 522
	}
	return nil, u.ErrNotFound
}

523
func (dht *IpfsDHT) GetLocal(key u.Key) ([]byte, error) {
524
	v, err := dht.datastore.Get(ds.NewKey(string(key)))
525 526 527 528 529 530 531 532 533
	if err != nil {
		return nil, err
	}
	return v.([]byte), nil
}

func (dht *IpfsDHT) PutLocal(key u.Key, value []byte) error {
	return dht.datastore.Put(ds.NewKey(string(key)), value)
}
534 535

func (dht *IpfsDHT) Update(p *peer.Peer) {
536 537 538 539 540 541 542 543 544 545 546 547 548 549 550
	for _, route := range dht.routes {
		removed := route.Update(p)
		// Only drop the connection if no tables refer to this peer
		if removed != nil {
			found := false
			for _, r := range dht.routes {
				if r.Find(removed.ID) != nil {
					found = true
					break
				}
			}
			if !found {
				dht.network.Drop(removed)
			}
		}
551 552
	}
}
Jeromy's avatar
Jeromy committed
553 554 555 556 557 558 559 560 561 562 563

// Look for a peer with a given ID connected to this dht
func (dht *IpfsDHT) Find(id peer.ID) (*peer.Peer, *kb.RoutingTable) {
	for _, table := range dht.routes {
		p := table.Find(id)
		if p != nil {
			return p, table
		}
	}
	return nil, nil
}
564 565 566 567 568 569 570 571 572 573

func (dht *IpfsDHT) findPeerSingle(p *peer.Peer, id peer.ID, timeout time.Duration, level int) (*PBDHTMessage, error) {
	pmes := DHTMessage{
		Type:  PBDHTMessage_FIND_NODE,
		Key:   string(id),
		Id:    GenerateMessageID(),
		Value: []byte{byte(level)},
	}

	mes := swarm.NewMessage(p, pmes.ToProtobuf())
574
	listenChan := dht.listener.Listen(pmes.Id, 1, time.Minute)
575
	t := time.Now()
576 577 578 579
	dht.network.Send(mes)
	after := time.After(timeout)
	select {
	case <-after:
580
		dht.listener.Unlisten(pmes.Id)
581 582
		return nil, u.ErrTimeout
	case resp := <-listenChan:
583 584
		roundtrip := time.Since(t)
		resp.Peer.SetLatency(roundtrip)
585 586 587 588 589 590 591 592 593
		pmes_out := new(PBDHTMessage)
		err := proto.Unmarshal(resp.Data, pmes_out)
		if err != nil {
			return nil, err
		}

		return pmes_out, nil
	}
}
594 595 596 597 598 599

func (dht *IpfsDHT) PrintTables() {
	for _, route := range dht.routes {
		route.Print()
	}
}
Jeromy's avatar
Jeromy committed
600 601 602 603 604 605 606 607 608 609 610

func (dht *IpfsDHT) findProvidersSingle(p *peer.Peer, key u.Key, level int, timeout time.Duration) (*PBDHTMessage, error) {
	pmes := DHTMessage{
		Type:  PBDHTMessage_GET_PROVIDERS,
		Key:   string(key),
		Id:    GenerateMessageID(),
		Value: []byte{byte(level)},
	}

	mes := swarm.NewMessage(p, pmes.ToProtobuf())

611
	listenChan := dht.listener.Listen(pmes.Id, 1, time.Minute)
Jeromy's avatar
Jeromy committed
612 613 614 615
	dht.network.Send(mes)
	after := time.After(timeout)
	select {
	case <-after:
616
		dht.listener.Unlisten(pmes.Id)
Jeromy's avatar
Jeromy committed
617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656
		return nil, u.ErrTimeout
	case resp := <-listenChan:
		u.DOut("FindProviders: got response.")
		pmes_out := new(PBDHTMessage)
		err := proto.Unmarshal(resp.Data, pmes_out)
		if err != nil {
			return nil, err
		}

		return pmes_out, nil
	}
}

func (dht *IpfsDHT) addPeerList(key u.Key, peers []*PBDHTMessage_PBPeer) []*peer.Peer {
	var prov_arr []*peer.Peer
	for _, prov := range peers {
		// Dont add outselves to the list
		if peer.ID(prov.GetId()).Equal(dht.self.ID) {
			continue
		}
		// Dont add someone who is already on the list
		p := dht.network.Find(u.Key(prov.GetId()))
		if p == nil {
			u.DOut("given provider %s was not in our network already.", peer.ID(prov.GetId()).Pretty())
			maddr, err := ma.NewMultiaddr(prov.GetAddr())
			if err != nil {
				u.PErr("error connecting to new peer: %s", err)
				continue
			}
			p, err = dht.network.GetConnection(peer.ID(prov.GetId()), maddr)
			if err != nil {
				u.PErr("error connecting to new peer: %s", err)
				continue
			}
		}
		dht.addProviderEntry(key, p)
		prov_arr = append(prov_arr, p)
	}
	return prov_arr
}