@cryptotaxi247 / kubo / commits / 171f96b79

update messages and add some new code around handling/creating messages

Jeromy Johnson committed Jul 29, 2014 at 14:50 UTC 171f96b794aae13cafe72d81705ba246574fb0aa
3 files changed +44 -5
routing/dht/dht.go
+32 -2
@@ -2,6 +2,7 @@ package dht
2
3 import (
4 swarm "github.com/jbenet/go-ipfs/swarm"
5 + "sync"
6 )
7
8 // TODO. SEE https://github.com/jbenet/node-ipfs/blob/master/submodules/ipfs-dht/index.js
@@ -10,13 +11,42 @@ import (
11 // IpfsDHT is an implementation of Kademlia with Coral and S/Kademlia modifications.
12 // It is used to implement the base IpfsRouting module.
13 type IpfsDHT struct {
13 - routes RoutingTable
14 + routes RoutingTable
15
15 - network *swarm.Swarm
16 + network *swarm.Swarm
17 +
18 + listeners map[uint64]chan swarm.Message
19 + listenLock sync.RWMutex
20 }
21
22 +// Read in all messages from swarm and handle them appropriately
23 +// NOTE: this function is just a quick sketch
24 func (dht *IpfsDHT) handleMessages() {
25 for mes := range dht.network.Chan.Incoming {
26 + for {
27 + select {
28 + case mes := <-dht.network.Chan.Incoming:
29 + // Unmarshal message
30 + dht.listenLock.RLock()
31 + ch, ok := dht.listeners[id]
32 + dht.listenLock.RUnlock()
33 + if ok {
34 + // Send message to waiting goroutine
35 + ch <- mes
36 + }
37
38 + //case closeChan: or something
39 + }
40 + }
41 }
42 }
43 +
44 +// Register a handler for a specific message ID, used for getting replies
45 +// to certain messages (i.e. response to a GET_VALUE message)
46 +func (dht *IpfsDHT) ListenFor(mesid uint64) <-chan swarm.Message {
47 + lchan := make(chan swarm.Message)
48 + dht.listenLock.Lock()
49 + dht.listeners[mesid] = lchan
50 + dht.listenLock.Unlock()
51 + return lchan
52 +}
routing/dht/messages.proto
+5 -2
@@ -6,11 +6,14 @@ message DHTMessage {
6 enum MessageType {
7 PUT_VALUE = 0;
8 GET_VALUE = 1;
9 - PING = 2;
10 - FIND_NODE = 3;
9 + ADD_PROVIDER = 2;
10 + GET_PROVIDERS = 3;
11 + FIND_NODE = 4;
12 + PING = 5;
13 }
14
15 required MessageType type = 1;
16 optional string key = 2;
17 optional bytes value = 3;
18 + required int64 id = 4;
19 }
routing/dht/routing.go
+7 -1
@@ -43,14 +43,20 @@ func (s *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
43 mes.Data = []byte(pmes.String())
44 mes.Peer = p
45
46 - response_chan := s.network.ListenFor(pmes.Id)
46 + response_chan := s.ListenFor(pmes.Id)
47
48 + // Wait for either the response or a timeout
49 timeup := time.After(timeout)
50 select {
51 case <-timeup:
52 + // TODO: unregister listener
53 return nil, timeoutError
54 case resp := <-response_chan:
55 + return resp.Data, nil
56 }
57 +
58 + // Should never be hit
59 + return nil, nil
60 }
61
62