@cryptotaxi247 / kubo / commits / dc451fba2

implement find peer rpc

Jeromy committed Aug 5, 2014 at 20:31 UTC dc451fba2d325d2e157e122f5d59e480f89c9131
4 files changed +67 -22
identify/identify.go
+1 -7
@@ -10,13 +10,7 @@ import (
10 // Perform initial communication with this peer to share node ID's and
11 // initiate communication
12 func Handshake(self, remote *peer.Peer, in, out chan []byte) error {
13 -
14 - // temporary:
15 - // put your own id in a 16byte buffer and send that over to
16 - // the peer as your ID, then wait for them to send their ID.
17 - // Once that trade is finished, the handshake is complete and
18 - // both sides should 'trust' each other
19 -
13 + // TODO: make this more... secure.
14 out <- self.ID
15 resp := <-in
16 remote.ID = peer.ID(resp)
routing/dht/dht.go
+51 -14
@@ -73,6 +73,7 @@ func (dht *IpfsDHT) Start() {
73 }
74
75 // Connect to a new peer at the given address
76 +// TODO: move this into swarm
77 func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
78 if addr == nil {
79 panic("addr was nil!")
@@ -90,9 +91,21 @@ func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
91 return nil, err
92 }
93
94 + // Send node an address that you can be reached on
95 + myaddr := dht.self.NetAddress("tcp")
96 + mastr,err := myaddr.String()
97 + if err != nil {
98 + panic("No local address to send")
99 + }
100 +
101 + conn.Outgoing.MsgChan <- []byte(mastr)
102 +
103 dht.network.StartConn(conn)
104
95 - dht.routes.Update(peer)
105 + removed := dht.routes.Update(peer)
106 + if removed != nil {
107 + panic("need to remove this peer.")
108 + }
109 return peer, nil
110 }
111
@@ -115,7 +128,10 @@ func (dht *IpfsDHT) handleMessages() {
128 }
129
130 // Update peers latest visit in routing table
118 - dht.routes.Update(mes.Peer)
131 + removed := dht.routes.Update(mes.Peer)
132 + if removed != nil {
133 + panic("Need to handle removed peer.")
134 + }
135
136 // Note: not sure if this is the correct place for this
137 if pmes.GetResponse() {
@@ -140,7 +156,7 @@ func (dht *IpfsDHT) handleMessages() {
156 case DHTMessage_PUT_VALUE:
157 dht.handlePutValue(mes.Peer, pmes)
158 case DHTMessage_FIND_NODE:
143 - dht.handleFindNode(mes.Peer, pmes)
159 + dht.handleFindPeer(mes.Peer, pmes)
160 case DHTMessage_ADD_PROVIDER:
161 dht.handleAddProvider(mes.Peer, pmes)
162 case DHTMessage_GET_PROVIDERS:
@@ -171,14 +187,14 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *DHTMessage) {
187 mes := swarm.NewMessage(p, resp.ToProtobuf())
188 dht.network.Chan.Outgoing <- mes
189 } else if err == ds.ErrNotFound {
174 - // Find closest node(s) to desired key and reply with that info
190 + // Find closest peer(s) to desired key and reply with that info
191 // TODO: this will need some other metadata in the protobuf message
176 - // to signal to the querying node that the data its receiving
177 - // is actually a list of other nodes
192 + // to signal to the querying peer that the data its receiving
193 + // is actually a list of other peer
194 }
195 }
196
181 -// Store a value in this nodes local storage
197 +// Store a value in this peer local storage
198 func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *DHTMessage) {
199 dskey := ds.NewKey(pmes.GetKey())
200 err := dht.datastore.Put(dskey, pmes.GetValue())
@@ -189,7 +205,7 @@ func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *DHTMessage) {
205 }
206
207 func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *DHTMessage) {
192 - resp := &pDHTMessage{
208 + resp := pDHTMessage{
209 Type: pmes.GetType(),
210 Response: true,
211 Id: pmes.GetId(),
@@ -198,8 +214,29 @@ func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *DHTMessage) {
214 dht.network.Chan.Outgoing <-swarm.NewMessage(p, resp.ToProtobuf())
215 }
216
201 -func (dht *IpfsDHT) handleFindNode(p *peer.Peer, pmes *DHTMessage) {
202 - panic("Not implemented.")
217 +func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *DHTMessage) {
218 + closest := dht.routes.NearestPeer(convertKey(u.Key(pmes.GetKey())))
219 + if closest == nil {
220 + }
221 +
222 + if len(closest.Addresses) == 0 {
223 + panic("no addresses for connected peer...")
224 + }
225 +
226 + addr,err := closest.Addresses[0].String()
227 + if err != nil {
228 + panic(err)
229 + }
230 +
231 + resp := pDHTMessage{
232 + Type: pmes.GetType(),
233 + Response: true,
234 + Id: pmes.GetId(),
235 + Value: []byte(addr),
236 + }
237 +
238 + mes := swarm.NewMessage(p, resp.ToProtobuf())
239 + dht.network.Chan.Outgoing <-mes
240 }
241
242 func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
@@ -269,13 +306,13 @@ func (dht *IpfsDHT) Unlisten(mesid uint64) {
306 close(ch)
307 }
308
272 -// Stop all communications from this node and shut down
309 +// Stop all communications from this peer and shut down
310 func (dht *IpfsDHT) Halt() {
311 dht.shutdown <- struct{}{}
312 dht.network.Close()
313 }
314
278 -// Ping a node, log the time it took
315 +// Ping a peer, log the time it took
316 func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
317 // Thoughts: maybe this should accept an ID and do a peer lookup?
318 u.DOut("Enter Ping.")
@@ -294,8 +331,8 @@ func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
331 u.POut("Ping took %s.", roundtrip.String())
332 return nil
333 case <-tout:
297 - // Timed out, think about removing node from network
298 - u.DOut("Ping node timed out.")
334 + // Timed out, think about removing peer from network
335 + u.DOut("Ping peer timed out.")
336 return u.ErrTimeout
337 }
338 }
routing/dht/routing.go
+7 -1
@@ -188,6 +188,12 @@ func (s *IpfsDHT) FindPeer(id peer.ID, timeout time.Duration) (*peer.Peer, error
188 if err != nil {
189 return nil, err
190 }
191 - panic("Not yet implemented.")
191 + addr := string(pmes_out.GetValue())
192 + maddr, err := ma.NewMultiaddr(addr)
193 + if err != nil {
194 + return nil, err
195 + }
196 +
197 + return s.Connect(maddr)
198 }
199 }
swarm/swarm.go
+8
@@ -163,6 +163,14 @@ func (s *Swarm) handleNewConn(nconn net.Conn) {
163 panic(err)
164 }
165
166 + // Get address to contact remote peer from
167 + addr := <-conn.Incoming.MsgChan
168 + maddr, err := ma.NewMultiaddr(string(addr))
169 + if err != nil {
170 + u.PErr("Got invalid address from peer.")
171 + }
172 + p.AddAddress(maddr)
173 +
174 s.StartConn(conn)
175 }
176