@cryptotaxi247 / kubo / commits / 61f13ea7f

begin planning of identification process

Jeromy Johnson committed Jul 31, 2014 at 17:43 UTC 61f13ea7f7ef41defc9cb65ad5a7e818fc1501af
4 files changed +138 -14
identify/identify.go new
+28
@@ -0,0 +1,28 @@
1 +// The identify package handles how peers identify with eachother upon
2 +// connection to the network
3 +package identify
4 +
5 +import (
6 + peer "github.com/jbenet/go-ipfs/peer"
7 + swarm "github.com/jbenet/go-ipfs/swarm"
8 +)
9 +
10 +// Perform initial communication with this peer to share node ID's and
11 +// initiate communication
12 +func Handshake(self *peer.Peer, conn *swarm.Conn) 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 +
20 + id := make([]byte, 16)
21 + copy(id, self.ID)
22 +
23 + conn.Outgoing.MsgChan <- id
24 + resp := <-conn.Incoming.MsgChan
25 + conn.Peer.ID = peer.ID(resp)
26 +
27 + return nil
28 +}
identify/message.proto new
+3
@@ -0,0 +1,3 @@
1 +message Identify {
2 + required bytes id = 1;
3 +}
routing/dht/dht.go
+66 -8
@@ -2,10 +2,14 @@ package dht
2
3 import (
4 "sync"
5 + "time"
6
7 peer "github.com/jbenet/go-ipfs/peer"
8 swarm "github.com/jbenet/go-ipfs/swarm"
9 u "github.com/jbenet/go-ipfs/util"
10 + identify "github.com/jbenet/go-ipfs/identify"
11 +
12 + ma "github.com/jbenet/go-multiaddr"
13
14 ds "github.com/jbenet/datastore.go"
15
@@ -35,15 +39,44 @@ type IpfsDHT struct {
39 shutdown chan struct{}
40 }
41
38 -func NewDHT(p *peer.Peer) *IpfsDHT {
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)
40 - dht.self = p
45 +
46 dht.network = swarm.NewSwarm(p)
47 + //TODO: should Listen return an error?
48 + dht.network.Listen()
49 +
50 + dht.datastore = ds.NewMapDatastore()
51 +
52 + dht.self = p
53 dht.listeners = make(map[uint64]chan *swarm.Message)
54 dht.shutdown = make(chan struct{})
44 - return dht
55 + return dht, nil
56 +}
57 +
58 +// Connect to a new peer at the given address
59 +func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) error {
60 + peer := new(peer.Peer)
61 + peer.AddAddress(addr)
62 +
63 + conn,err := swarm.Dial("tcp", peer)
64 + if err != nil {
65 + return err
66 + }
67 +
68 + err = identify.Handshake(dht.self, conn)
69 + if err != nil {
70 + return err
71 + }
72 +
73 + dht.network.StartConn(conn.Peer.Key(), conn)
74 +
75 + // TODO: Add this peer to our routing table
76 + return nil
77 }
78
79 +
80 // Read in all messages from swarm and handle them appropriately
81 // NOTE: this function is just a quick sketch
82 func (dht *IpfsDHT) handleMessages() {
@@ -134,11 +167,9 @@ func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *DHTMessage) {
167 resp := new(DHTMessage)
168 resp.Id = pmes.Id
169 resp.Response = &isResponse
170 + resp.Type = pmes.Type
171
138 - mes := new(swarm.Message)
139 - mes.Peer = p
140 - mes.Data = []byte(resp.String())
141 - dht.network.Chan.Outgoing <- mes
172 + dht.network.Chan.Outgoing <-swarm.NewMessage(p, []byte(resp.String()))
173 }
174
175
@@ -162,9 +193,36 @@ func (dht *IpfsDHT) Unlisten(mesid uint64) {
193 close(ch)
194 }
195
165 -
196 // Stop all communications from this node and shut down
197 func (dht *IpfsDHT) Halt() {
198 dht.shutdown <- struct{}{}
199 dht.network.Close()
200 }
201 +
202 +// Ping a node, log the time it took
203 +func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) {
204 + // Thoughts: maybe this should accept an ID and do a peer lookup?
205 + id := GenerateMessageID()
206 + mes_type := DHTMessage_PING
207 + pmes := new(DHTMessage)
208 + pmes.Id = &id
209 + pmes.Type = &mes_type
210 +
211 + mes := new(swarm.Message)
212 + mes.Peer = p
213 + mes.Data = []byte(pmes.String())
214 +
215 + before := time.Now()
216 + response_chan := dht.ListenFor(id)
217 + dht.network.Chan.Outgoing <- mes
218 +
219 + tout := time.After(timeout)
220 + select {
221 + case <-response_chan:
222 + roundtrip := time.Since(before)
223 + u.DOut("Ping took %s.", roundtrip.String())
224 + case <-tout:
225 + // Timed out, think about removing node from network
226 + u.DOut("Ping node timed out.")
227 + }
228 +}
swarm/swarm.go
+41 -6
@@ -2,11 +2,12 @@ package swarm
2
3 import (
4 "fmt"
5 + "net"
6 + "sync"
7 +
8 peer "github.com/jbenet/go-ipfs/peer"
9 u "github.com/jbenet/go-ipfs/util"
10 ma "github.com/jbenet/go-multiaddr"
8 - "net"
9 - "sync"
11 )
12
13 // Message represents a packet of information sent to or received from a
@@ -19,6 +20,14 @@ type Message struct {
20 Data []byte
21 }
22
23 +// Cleaner looking helper function to make a new message struct
24 +func NewMessage(p *peer.Peer, data []byte) *Message {
25 + return &Message{
26 + Peer: p,
27 + Data: data,
28 + }
29 +}
30 +
31 // Chan is a swam channel, which provides duplex communication and errors.
32 type Chan struct {
33 Outgoing chan *Message
@@ -87,7 +96,8 @@ func (s *Swarm) connListen(maddr *ma.Multiaddr) error {
96 for {
97 nconn, err := list.Accept()
98 if err != nil {
90 - u.PErr("Failed to accept connection: %s - %s", netstr, addr)
99 + u.PErr("Failed to accept connection: %s - %s [%s]", netstr,
100 + addr, err)
101 return
102 }
103 go s.handleNewConn(nconn)
@@ -99,7 +109,27 @@ func (s *Swarm) connListen(maddr *ma.Multiaddr) error {
109
110 // Handle getting ID from this peer and adding it into the map
111 func (s *Swarm) handleNewConn(nconn net.Conn) {
102 - panic("Not yet implemented!")
112 + p := MakePeerFromConn(nconn)
113 +
114 + var addr *ma.Multiaddr
115 +
116 + //naddr := nconn.RemoteAddr()
117 + //addr := ma.FromDialArgs(naddr.Network(), naddr.String())
118 +
119 + conn := &Conn{
120 + Peer: p,
121 + Addr: addr,
122 + Conn: nconn,
123 + }
124 +
125 + newConnChans(conn)
126 + go s.fanIn(conn)
127 +}
128 +
129 +// Negotiate with peer for its ID and create a peer object
130 +// TODO: this might belong in the peer package
131 +func MakePeerFromConn(conn net.Conn) *peer.Peer {
132 + panic("Not yet implemented.")
133 }
134
135 // Close closes a swarm.
@@ -140,6 +170,11 @@ func (s *Swarm) Dial(peer *peer.Peer) (*Conn, error) {
170 return nil, err
171 }
172
173 + s.StartConn(k, conn)
174 + return conn, nil
175 +}
176 +
177 +func (s *Swarm) StartConn(k u.Key, conn *Conn) {
178 // add to conns
179 s.connsLock.Lock()
180 s.conns[k] = conn
@@ -147,7 +182,6 @@ func (s *Swarm) Dial(peer *peer.Peer) (*Conn, error) {
182
183 // kick off reader goroutine
184 go s.fanIn(conn)
150 - return conn, nil
185 }
186
187 // Handles the unwrapping + sending of messages to the right connection.
@@ -165,7 +199,8 @@ func (s *Swarm) fanOut() {
199 conn, found := s.conns[msg.Peer.Key()]
200 s.connsLock.RUnlock()
201 if !found {
168 - e := fmt.Errorf("Sent msg to peer without open conn: %v", msg.Peer)
202 + e := fmt.Errorf("Sent msg to peer without open conn: %v",
203 + msg.Peer)
204 s.Chan.Errors <- e
205 }
206