@cryptotaxi247 / kubo / commits / 71c7c5844

providers interface is coming along nicely

Jeromy committed Aug 5, 2014 at 18:32 UTC 71c7c5844a5274b55243277690241d255b93f1af
7 files changed +80 -35
identify/identify.go
+1 -1
@@ -20,7 +20,7 @@ func Handshake(self, remote *peer.Peer, in, out chan []byte) error {
20 out <- self.ID
21 resp := <-in
22 remote.ID = peer.ID(resp)
23 - u.DOut("Got node id: %s", string(remote.ID))
23 + u.DOut("identify: Got node id: %s", remote.ID.Pretty())
24
25 return nil
26 }
peer/peer.go
+6
@@ -1,6 +1,8 @@
1 package peer
2
3 import (
4 + "encoding/hex"
5 +
6 u "github.com/jbenet/go-ipfs/util"
7 ma "github.com/jbenet/go-multiaddr"
8 mh "github.com/jbenet/go-multihash"
@@ -16,6 +18,10 @@ func (id ID) Equal(other ID) bool {
18 return bytes.Equal(id, other)
19 }
20
21 +func (id ID) Pretty() string {
22 + return hex.EncodeToString(id)
23 +}
24 +
25 // Map maps Key (string) : *Peer (slices are not comparable).
26 type Map map[u.Key]*Peer
27
routing/dht/dht.go
+28 -11
@@ -61,6 +61,7 @@ func NewDHT(p *peer.Peer) (*IpfsDHT, error) {
61 dht.datastore = ds.NewMapDatastore()
62 dht.self = p
63 dht.listeners = make(map[uint64]chan *swarm.Message)
64 + dht.providers = make(map[u.Key][]*peer.Peer)
65 dht.shutdown = make(chan struct{})
66 dht.routes = NewRoutingTable(20, convertPeerID(p.ID))
67 return dht, nil
@@ -101,9 +102,11 @@ func (dht *IpfsDHT) handleMessages() {
102 u.DOut("Begin message handling routine")
103 for {
104 select {
104 - case mes := <-dht.network.Chan.Incoming:
105 - u.DOut("recieved message from swarm.")
106 -
105 + case mes,ok := <-dht.network.Chan.Incoming:
106 + if !ok {
107 + u.DOut("handleMessages closing, bad recv on incoming")
108 + return
109 + }
110 pmes := new(DHTMessage)
111 err := proto.Unmarshal(mes.Data, pmes)
112 if err != nil {
@@ -121,15 +124,16 @@ func (dht *IpfsDHT) handleMessages() {
124 dht.listenLock.RUnlock()
125 if ok {
126 ch <- mes
127 + } else {
128 + // this is expected behaviour during a timeout
129 + u.DOut("Received response with nobody listening...")
130 }
131
126 - // this is expected behaviour during a timeout
127 - u.DOut("Received response with nobody listening...")
132 continue
133 }
134 //
135
132 - u.DOut("Got message type: %d", pmes.GetType())
136 + u.DOut("Got message type: '%s' [id = %x]", mesNames[pmes.GetType()], pmes.GetId())
137 switch pmes.GetType() {
138 case DHTMessage_GET_VALUE:
139 dht.handleGetValue(mes.Peer, pmes)
@@ -138,13 +142,15 @@ func (dht *IpfsDHT) handleMessages() {
142 case DHTMessage_FIND_NODE:
143 dht.handleFindNode(mes.Peer, pmes)
144 case DHTMessage_ADD_PROVIDER:
145 + dht.handleAddProvider(mes.Peer, pmes)
146 case DHTMessage_GET_PROVIDERS:
147 + dht.handleGetProviders(mes.Peer, pmes)
148 case DHTMessage_PING:
149 dht.handlePing(mes.Peer, pmes)
150 }
151
152 case err := <-dht.network.Chan.Errors:
147 - panic(err)
153 + u.DErr("dht err: %s", err)
154 case <-dht.shutdown:
155 return
156 }
@@ -197,13 +203,16 @@ func (dht *IpfsDHT) handleFindNode(p *peer.Peer, pmes *DHTMessage) {
203 }
204
205 func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
206 + dht.providerLock.RLock()
207 providers := dht.providers[u.Key(pmes.GetKey())]
208 + dht.providerLock.RUnlock()
209 if providers == nil || len(providers) == 0 {
210 // ?????
211 + u.DOut("No known providers for requested key.")
212 }
213
214 // This is just a quick hack, formalize method of sending addrs later
206 - var addrs []string
215 + addrs := make(map[u.Key]string)
216 for _,prov := range providers {
217 ma := prov.NetAddress("tcp")
218 str,err := ma.String()
@@ -212,7 +221,7 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
221 continue
222 }
223
215 - addrs = append(addrs, str)
224 + addrs[prov.Key()] = str
225 }
226
227 data,err := json.Marshal(addrs)
@@ -225,6 +234,7 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
234 Key: pmes.GetKey(),
235 Value: data,
236 Id: pmes.GetId(),
237 + Response: true,
238 }
239
240 mes := swarm.NewMessage(p, resp.ToProtobuf())
@@ -234,8 +244,7 @@ func (dht *IpfsDHT) handleGetProviders(p *peer.Peer, pmes *DHTMessage) {
244 func (dht *IpfsDHT) handleAddProvider(p *peer.Peer, pmes *DHTMessage) {
245 //TODO: need to implement TTLs on providers
246 key := u.Key(pmes.GetKey())
237 - parr := dht.providers[key]
238 - dht.providers[key] = append(parr, p)
247 + dht.addProviderEntry(key, p)
248 }
249
250
@@ -290,3 +299,11 @@ func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
299 return u.ErrTimeout
300 }
301 }
302 +
303 +func (dht *IpfsDHT) addProviderEntry(key u.Key, p *peer.Peer) {
304 + u.DOut("Adding %s as provider for '%s'", p.Key().Pretty(), key)
305 + dht.providerLock.Lock()
306 + provs := dht.providers[key]
307 + dht.providers[key] = append(provs, p)
308 + dht.providerLock.Unlock()
309 +}
routing/dht/pDHTMessage.go
+11
@@ -9,6 +9,17 @@ type pDHTMessage struct {
9 Id uint64
10 }
11
12 +var mesNames [10]string
13 +
14 +func init() {
15 + mesNames[DHTMessage_ADD_PROVIDER] = "add provider"
16 + mesNames[DHTMessage_FIND_NODE] = "find node"
17 + mesNames[DHTMessage_GET_PROVIDERS] = "get providers"
18 + mesNames[DHTMessage_GET_VALUE] = "get value"
19 + mesNames[DHTMessage_PUT_VALUE] = "put value"
20 + mesNames[DHTMessage_PING] = "ping"
21 +}
22 +
23 func (m *pDHTMessage) ToProtobuf() *DHTMessage {
24 pmes := new(DHTMessage)
25 if m.Value != nil {
routing/dht/routing.go
+14 -10
@@ -19,7 +19,8 @@ var PoolSize = 6
19
20 // TODO: determine a way of creating and managing message IDs
21 func GenerateMessageID() uint64 {
22 - return uint64(rand.Uint32()) << 32 & uint64(rand.Uint32())
22 + //return (uint64(rand.Uint32()) << 32) & uint64(rand.Uint32())
23 + return uint64(rand.Uint32())
24 }
25
26 // This file implements the Routing interface for the IpfsDHT struct.
@@ -116,6 +117,7 @@ func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer,
117 mes := swarm.NewMessage(p, pmes.ToProtobuf())
118
119 listen_chan := s.ListenFor(pmes.Id)
120 + u.DOut("Find providers for: '%s'", key)
121 s.network.Chan.Outgoing <-mes
122 after := time.After(timeout)
123 select {
@@ -123,37 +125,39 @@ func (s *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Peer,
125 s.Unlisten(pmes.Id)
126 return nil, u.ErrTimeout
127 case resp := <-listen_chan:
128 + u.DOut("FindProviders: got response.")
129 pmes_out := new(DHTMessage)
130 err := proto.Unmarshal(resp.Data, pmes_out)
131 if err != nil {
132 return nil, err
133 }
131 - var addrs map[string]string
132 - err := json.Unmarshal(pmes_out.GetValue(), &addrs)
134 + var addrs map[u.Key]string
135 + err = json.Unmarshal(pmes_out.GetValue(), &addrs)
136 if err != nil {
137 return nil, err
138 }
139
137 - for key,addr := range addrs {
138 - p := s.network.Find(u.Key(key))
140 + var prov_arr []*peer.Peer
141 + for pid,addr := range addrs {
142 + p := s.network.Find(pid)
143 if p == nil {
144 maddr,err := ma.NewMultiaddr(addr)
145 if err != nil {
146 u.PErr("error connecting to new peer: %s", err)
147 continue
148 }
145 - p, err := s.Connect(maddr)
149 + p, err = s.Connect(maddr)
150 if err != nil {
151 u.PErr("error connecting to new peer: %s", err)
152 continue
153 }
154 }
151 - s.providerLock.Lock()
152 - prov_arr := s.providers[key]
153 - s.providers[key] = append(prov_arr, p)
154 - s.providerLock.Unlock()
155 + s.addProviderEntry(key, p)
156 + prov_arr = append(prov_arr, p)
157 }
158
159 + return prov_arr, nil
160 +
161 }
162 }
163
swarm/swarm.go
+15 -13
@@ -47,7 +47,7 @@ func NewChan(bufsize int) *Chan {
47 return &Chan{
48 Outgoing: make(chan *Message, bufsize),
49 Incoming: make(chan *Message, bufsize),
50 - Errors: make(chan error),
50 + Errors: make(chan error, bufsize),
51 Close: make(chan bool, bufsize),
52 }
53 }
@@ -137,7 +137,7 @@ func (s *Swarm) connListen(maddr *ma.Multiaddr) error {
137 if err != nil {
138 e := fmt.Errorf("Failed to accept connection: %s - %s [%s]",
139 netstr, addr, err)
140 - s.Chan.Errors <- e
140 + go func() {s.Chan.Errors <- e}()
141 return
142 }
143 go s.handleNewConn(nconn)
@@ -217,7 +217,7 @@ func (s *Swarm) StartConn(conn *Conn) {
217 panic("tried to start nil Conn!")
218 }
219
220 - u.DOut("Starting connection: %s", string(conn.Peer.ID))
220 + u.DOut("Starting connection: %s", conn.Peer.Key().Pretty())
221 // add to conns
222 s.connsLock.Lock()
223 s.conns[conn.Peer.Key()] = conn
@@ -234,10 +234,10 @@ func (s *Swarm) fanOut() {
234 case <-s.Chan.Close:
235 return // told to close.
236 case msg, ok := <-s.Chan.Outgoing:
237 - u.DOut("fanOut: outgoing message for: '%s'", msg.Peer.Key())
237 if !ok {
238 return
239 }
240 + u.DOut("fanOut: outgoing message for: '%s'", msg.Peer.Key().Pretty())
241
242 s.connsLock.RLock()
243 conn, found := s.conns[msg.Peer.Key()]
@@ -252,7 +252,6 @@ func (s *Swarm) fanOut() {
252
253 // queue it in the connection's buffer
254 conn.Outgoing.MsgChan <- msg.Data
255 - u.DOut("fanOut: message off.")
255 }
256 }
257 }
@@ -260,36 +259,39 @@ func (s *Swarm) fanOut() {
259 // Handles the receiving + wrapping of messages, per conn.
260 // Consider using reflect.Select with one goroutine instead of n.
261 func (s *Swarm) fanIn(conn *Conn) {
263 -Loop:
262 for {
263 select {
264 case <-s.Chan.Close:
265 // close Conn.
266 conn.Close()
269 - break Loop
267 + goto out
268
269 case <-conn.Closed:
272 - break Loop
270 + goto out
271
272 case data, ok := <-conn.Incoming.MsgChan:
275 - u.DOut("fanIn: got message from incoming channel.")
273 if !ok {
277 - e := fmt.Errorf("Error retrieving from conn: %v", conn)
274 + e := fmt.Errorf("Error retrieving from conn: %v", conn.Peer.Key().Pretty())
275 s.Chan.Errors <- e
279 - break Loop
276 + goto out
277 }
278
279 // wrap it for consumers.
280 msg := &Message{Peer: conn.Peer, Data: data}
281 s.Chan.Incoming <- msg
285 - u.DOut("fanIn: message off.")
282 }
283 }
284 +out:
285
286 s.connsLock.Lock()
287 delete(s.conns, conn.Peer.Key())
288 s.connsLock.Unlock()
289 }
290
294 -func (s *Swarm) Find(addr *ma.Multiaddr) {
291 +func (s *Swarm) Find(key u.Key) *peer.Peer {
292 + conn, found := s.conns[key]
293 + if !found {
294 + return nil
295 + }
296 + return conn.Peer
297 }
util/util.go
+5
@@ -6,6 +6,7 @@ import (
6 "os"
7 "os/user"
8 "strings"
9 + "encoding/hex"
10 )
11
12 // Debug is a global flag for debugging.
@@ -20,6 +21,10 @@ var ErrTimeout = fmt.Errorf("Error: Call timed out.")
21 // Key is a string representation of multihash for use with maps.
22 type Key string
23
24 +func (k Key) Pretty() string {
25 + return hex.EncodeToString([]byte(k))
26 +}
27 +
28 // Hash is the global IPFS hash function. uses multihash SHA2_256, 256 bits
29 func Hash(data []byte) (mh.Multihash, error) {
30 return mh.Sum(data, mh.SHA2_256, -1)