@cryptotaxi247 / kubo / commits / 92fb51d9a

finish basic communcations between nodes and add a test of the ping operation

Jeromy Johnson committed Aug 1, 2014 at 13:21 UTC 92fb51d9a2fef50c55d9d674a1813f68f77a6f4f
4 files changed +150 -41
routing/dht/dht.go
+33 -37
@@ -41,20 +41,22 @@ type IpfsDHT struct {
41
42 // Create a new DHT object with the given peer as the 'local' host
43 func NewDHT(p *peer.Peer) (*IpfsDHT, error) {
44 - dht := new(IpfsDHT)
45 -
46 - dht.network = swarm.NewSwarm(p)
47 - //TODO: should Listen return an error?
48 - dht.network.Listen()
44 + network := swarm.NewSwarm(p)
45 + err := network.Listen()
46 + if err != nil {
47 + return nil,err
48 + }
49
50 + dht := new(IpfsDHT)
51 + dht.network = network
52 dht.datastore = ds.NewMapDatastore()
51 -
53 dht.self = p
54 dht.listeners = make(map[uint64]chan *swarm.Message)
55 dht.shutdown = make(chan struct{})
56 return dht, nil
57 }
58
59 +// Start up background goroutines needed by the DHT
60 func (dht *IpfsDHT) Start() {
61 go dht.handleMessages()
62 }
@@ -111,6 +113,7 @@ func (dht *IpfsDHT) handleMessages() {
113 }
114 //
115
116 + u.DOut("Got message type: %d", pmes.GetType())
117 switch pmes.GetType() {
118 case DHTMessage_GET_VALUE:
119 dht.handleGetValue(mes.Peer, pmes)
@@ -121,7 +124,7 @@ func (dht *IpfsDHT) handleMessages() {
124 case DHTMessage_ADD_PROVIDER:
125 case DHTMessage_GET_PROVIDERS:
126 case DHTMessage_PING:
124 - dht.handleFindNode(mes.Peer, pmes)
127 + dht.handlePing(mes.Peer, pmes)
128 }
129
130 case err := <-dht.network.Chan.Errors:
@@ -136,18 +139,15 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *DHTMessage) {
139 dskey := ds.NewKey(pmes.GetKey())
140 i_val, err := dht.datastore.Get(dskey)
141 if err == nil {
139 - isResponse := true
140 - resp := new(DHTMessage)
141 - resp.Response = &isResponse
142 - resp.Id = pmes.Id
143 - resp.Key = pmes.Key
144 -
145 - val := i_val.([]byte)
146 - resp.Value = val
147 -
148 - mes := new(swarm.Message)
149 - mes.Peer = p
150 - mes.Data = []byte(resp.String())
142 + resp := &pDHTMessage{
143 + Response: true,
144 + Id: *pmes.Id,
145 + Key: *pmes.Key,
146 + Value: i_val.([]byte),
147 + }
148 +
149 + mes := swarm.NewMessage(p, resp.ToProtobuf())
150 + dht.network.Chan.Outgoing <- mes
151 } else if err == ds.ErrNotFound {
152 // Find closest node(s) to desired key and reply with that info
153 // TODO: this will need some other metadata in the protobuf message
@@ -167,13 +167,13 @@ func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *DHTMessage) {
167 }
168
169 func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *DHTMessage) {
170 - isResponse := true
171 - resp := new(DHTMessage)
172 - resp.Id = pmes.Id
173 - resp.Response = &isResponse
174 - resp.Type = pmes.Type
170 + resp := &pDHTMessage{
171 + Type: pmes.GetType(),
172 + Response: true,
173 + Id: pmes.GetId(),
174 + }
175
176 - dht.network.Chan.Outgoing <-swarm.NewMessage(p, []byte(resp.String()))
176 + dht.network.Chan.Outgoing <-swarm.NewMessage(p, resp.ToProtobuf())
177 }
178
179 func (dht *IpfsDHT) handleFindNode(p *peer.Peer, pmes *DHTMessage) {
@@ -199,6 +199,7 @@ func (dht *IpfsDHT) ListenFor(mesid uint64) <-chan *swarm.Message {
199 return lchan
200 }
201
202 +// Unregister the given message id from the listener map
203 func (dht *IpfsDHT) Unlisten(mesid uint64) {
204 dht.listenLock.Lock()
205 ch, ok := dht.listeners[mesid]
@@ -216,31 +217,26 @@ func (dht *IpfsDHT) Halt() {
217 }
218
219 // Ping a node, log the time it took
219 -func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) {
220 +func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
221 // Thoughts: maybe this should accept an ID and do a peer lookup?
222 u.DOut("Enter Ping.")
223
223 - id := GenerateMessageID()
224 - mes_type := DHTMessage_PING
225 - pmes := new(DHTMessage)
226 - pmes.Id = &id
227 - pmes.Type = &mes_type
228 -
229 - mes := new(swarm.Message)
230 - mes.Peer = p
231 - mes.Data = []byte(pmes.String())
224 + pmes := pDHTMessage{Id: GenerateMessageID(), Type: DHTMessage_PING}
225 + mes := swarm.NewMessage(p, pmes.ToProtobuf())
226
227 before := time.Now()
234 - response_chan := dht.ListenFor(id)
228 + response_chan := dht.ListenFor(pmes.Id)
229 dht.network.Chan.Outgoing <- mes
230
231 tout := time.After(timeout)
232 select {
233 case <-response_chan:
234 roundtrip := time.Since(before)
241 - u.DOut("Ping took %s.", roundtrip.String())
235 + u.POut("Ping took %s.", roundtrip.String())
236 + return nil
237 case <-tout:
238 // Timed out, think about removing node from network
239 u.DOut("Ping node timed out.")
240 + return u.ErrTimeout
241 }
242 }
routing/dht/dht_test.go new
+55
@@ -0,0 +1,55 @@
1 +package dht
2 +
3 +import (
4 + "testing"
5 + peer "github.com/jbenet/go-ipfs/peer"
6 + ma "github.com/jbenet/go-multiaddr"
7 + u "github.com/jbenet/go-ipfs/util"
8 +
9 + "time"
10 +)
11 +
12 +func TestPing(t *testing.T) {
13 + u.Debug = false
14 + addr_a,err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1234")
15 + if err != nil {
16 + t.Fatal(err)
17 + }
18 + addr_b,err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/5678")
19 + if err != nil {
20 + t.Fatal(err)
21 + }
22 +
23 + peer_a := new(peer.Peer)
24 + peer_a.AddAddress(addr_a)
25 + peer_a.ID = peer.ID([]byte("peer_a"))
26 +
27 + peer_b := new(peer.Peer)
28 + peer_b.AddAddress(addr_b)
29 + peer_b.ID = peer.ID([]byte("peer_b"))
30 +
31 + dht_a,err := NewDHT(peer_a)
32 + if err != nil {
33 + t.Fatal(err)
34 + }
35 +
36 + dht_b,err := NewDHT(peer_b)
37 + if err != nil {
38 + t.Fatal(err)
39 + }
40 +
41 +
42 + dht_a.Start()
43 + dht_b.Start()
44 +
45 + err = dht_a.Connect(addr_b)
46 + if err != nil {
47 + t.Fatal(err)
48 + }
49 +
50 + //Test that we can ping the node
51 + err = dht_a.Ping(peer_b, time.Second * 2)
52 + if err != nil {
53 + t.Fatal(err)
54 + }
55 +}
routing/dht/pDHTMessage.go new
+24
@@ -0,0 +1,24 @@
1 +package dht
2 +
3 +// A helper struct to make working with protbuf types easier
4 +type pDHTMessage struct {
5 + Type DHTMessage_MessageType
6 + Key string
7 + Value []byte
8 + Response bool
9 + Id uint64
10 +}
11 +
12 +func (m *pDHTMessage) ToProtobuf() *DHTMessage {
13 + pmes := new(DHTMessage)
14 + if m.Value != nil {
15 + pmes.Value = m.Value
16 + }
17 +
18 + pmes.Type = &m.Type
19 + pmes.Key = &m.Key
20 + pmes.Response = &m.Response
21 + pmes.Id = &m.Id
22 +
23 + return pmes
24 +}
swarm/swarm.go
+38 -4
@@ -9,6 +9,7 @@ import (
9 u "github.com/jbenet/go-ipfs/util"
10 ma "github.com/jbenet/go-multiaddr"
11 ident "github.com/jbenet/go-ipfs/identify"
12 + proto "code.google.com/p/goprotobuf/proto"
13 )
14
15 // Message represents a packet of information sent to or received from a
@@ -22,10 +23,14 @@ type Message struct {
23 }
24
25 // Cleaner looking helper function to make a new message struct
25 -func NewMessage(p *peer.Peer, data []byte) *Message {
26 +func NewMessage(p *peer.Peer, data proto.Message) *Message {
27 + bytes,err := proto.Marshal(data)
28 + if err != nil {
29 + panic(err)
30 + }
31 return &Message{
32 Peer: p,
28 - Data: data,
33 + Data: bytes,
34 }
35 }
36
@@ -47,6 +52,25 @@ func NewChan(bufsize int) *Chan {
52 }
53 }
54
55 +// Contains a set of errors mapping to each of the swarms addresses
56 +// that were listened on
57 +type SwarmListenErr struct {
58 + Errors []error
59 +}
60 +
61 +func (se *SwarmListenErr) Error() string {
62 + if se == nil {
63 + return "<nil error>"
64 + }
65 + var out string
66 + for i,v := range se.Errors {
67 + if v != nil {
68 + out += fmt.Sprintf("%d: %s\n", i, v)
69 + }
70 + }
71 + return out
72 +}
73 +
74 // Swarm is a connection muxer, allowing connections to other peers to
75 // be opened and closed, while still using the same Chan for all
76 // communication. The Chan sends/receives Messages, which note the
@@ -71,13 +95,23 @@ func NewSwarm(local *peer.Peer) *Swarm {
95 }
96
97 // Open listeners for each network the swarm should listen on
74 -func (s *Swarm) Listen() {
75 - for _, addr := range s.local.Addresses {
98 +func (s *Swarm) Listen() error {
99 + var ret_err *SwarmListenErr
100 + for i, addr := range s.local.Addresses {
101 err := s.connListen(addr)
102 if err != nil {
103 + if ret_err != nil {
104 + ret_err = new(SwarmListenErr)
105 + ret_err.Errors = make([]error, len(s.local.Addresses))
106 + }
107 + ret_err.Errors[i] = err
108 u.PErr("Failed to listen on: %s [%s]", addr, err)
109 }
110 }
111 + if ret_err == nil {
112 + return nil
113 + }
114 + return ret_err
115 }
116
117 // Listen for new connections on the given multiaddr