@cryptotaxi247 / kubo / commits / 3c1580e1d

fix a few race conditions and add in newlines to print statements

Jeromy committed Aug 17, 2014 at 20:17 UTC 3c1580e1d10629c8a50783d9cdc1af4dbf79d747
5 files changed +45 -30
identify/identify.go
+1 -1
@@ -14,7 +14,7 @@ func Handshake(self, remote *peer.Peer, in, out chan []byte) error {
14 out <- self.ID
15 resp := <-in
16 remote.ID = peer.ID(resp)
17 - u.DOut("[%s] identify: Got node id: %s", self.ID.Pretty(), remote.ID.Pretty())
17 + u.DOut("[%s] identify: Got node id: %s\n", self.ID.Pretty(), remote.ID.Pretty())
18
19 return nil
20 }
routing/dht/dht.go
+26 -19
@@ -34,6 +34,7 @@ type IpfsDHT struct {
34
35 // Local data
36 datastore ds.Datastore
37 + dslock sync.Mutex
38
39 // Map keys to peers that can provide their value
40 providers map[u.Key][]*providerInfo
@@ -78,7 +79,7 @@ func (dht *IpfsDHT) Start() {
79 // Connect to a new peer at the given address, ping and add to the routing table
80 func (dht *IpfsDHT) Connect(addr *ma.Multiaddr) (*peer.Peer, error) {
81 maddrstr, _ := addr.String()
81 - u.DOut("Connect to new peer: %s", maddrstr)
82 + u.DOut("Connect to new peer: %s\n", maddrstr)
83 npeer, err := dht.network.ConnectNew(addr)
84 if err != nil {
85 return nil, err
@@ -113,7 +114,7 @@ func (dht *IpfsDHT) handleMessages() {
114 pmes := new(PBDHTMessage)
115 err := proto.Unmarshal(mes.Data, pmes)
116 if err != nil {
116 - u.PErr("Failed to decode protobuf message: %s", err)
117 + u.PErr("Failed to decode protobuf message: %s\n", err)
118 continue
119 }
120
@@ -126,7 +127,7 @@ func (dht *IpfsDHT) handleMessages() {
127 }
128 //
129
129 - u.DOut("[peer: %s]\nGot message type: '%s' [id = %x, from = %s]",
130 + u.DOut("[peer: %s]\nGot message type: '%s' [id = %x, from = %s]\n",
131 dht.self.ID.Pretty(),
132 PBDHTMessage_MessageType_name[int32(pmes.GetType())],
133 pmes.GetId(), mes.Peer.ID.Pretty())
@@ -148,7 +149,7 @@ func (dht *IpfsDHT) handleMessages() {
149 }
150
151 case err := <-ch.Errors:
151 - u.PErr("dht err: %s", err)
152 + u.PErr("dht err: %s\n", err)
153 case <-dht.shutdown:
154 checkTimeouts.Stop()
155 return
@@ -187,7 +188,7 @@ func (dht *IpfsDHT) putValueToNetwork(p *peer.Peer, key string, value []byte) er
188 }
189
190 func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
190 - u.DOut("handleGetValue for key: %s", pmes.GetKey())
191 + u.DOut("handleGetValue for key: %s\n", pmes.GetKey())
192 dskey := ds.NewKey(pmes.GetKey())
193 resp := &Message{
194 Response: true,
@@ -201,9 +202,11 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
202 resp.Value = iVal.([]byte)
203 } else if err == ds.ErrNotFound {
204 // Check if we know any providers for the requested value
205 + dht.providerLock.RLock()
206 provs, ok := dht.providers[u.Key(pmes.GetKey())]
207 + dht.providerLock.RUnlock()
208 if ok && len(provs) > 0 {
206 - u.DOut("handleGetValue returning %d provider[s]", len(provs))
209 + u.DOut("handleGetValue returning %d provider[s]\n", len(provs))
210 for _, prov := range provs {
211 resp.Peers = append(resp.Peers, prov.Value)
212 }
@@ -219,7 +222,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
222 } else {
223 level = int(pmes.GetValue()[0]) // Using value field to specify cluster level
224 }
222 - u.DOut("handleGetValue searching level %d clusters", level)
225 + u.DOut("handleGetValue searching level %d clusters\n", level)
226
227 closer := dht.routingTables[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
228
@@ -233,7 +236,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *PBDHTMessage) {
236 resp.Peers = nil
237 u.DOut("handleGetValue could not find a closer node than myself.")
238 } else {
236 - u.DOut("handleGetValue returning a closer peer: '%s'", closer.ID.Pretty())
239 + u.DOut("handleGetValue returning a closer peer: '%s'\n", closer.ID.Pretty())
240 resp.Peers = []*peer.Peer{closer}
241 }
242 }
@@ -249,6 +252,8 @@ out:
252
253 // Store a value in this peer local storage
254 func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *PBDHTMessage) {
255 + dht.dslock.Lock()
256 + defer dht.dslock.Unlock()
257 dskey := ds.NewKey(pmes.GetKey())
258 err := dht.datastore.Put(dskey, pmes.GetValue())
259 if err != nil {
@@ -278,7 +283,7 @@ func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *PBDHTMessage) {
283 dht.network.Send(mes)
284 }()
285 level := pmes.GetValue()[0]
281 - u.DOut("handleFindPeer: searching for '%s'", peer.ID(pmes.GetKey()).Pretty())
286 + u.DOut("handleFindPeer: searching for '%s'\n", peer.ID(pmes.GetKey()).Pretty())
287 closest := dht.routingTables[level].NearestPeer(kb.ConvertKey(u.Key(pmes.GetKey())))
288 if closest == nil {
289 u.PErr("handleFindPeer: could not find anything.")
@@ -295,7 +300,7 @@ func (dht *IpfsDHT) handleFindPeer(p *peer.Peer, pmes *PBDHTMessage) {
300 return
301 }
302
298 - u.DOut("handleFindPeer: sending back '%s'", closest.ID.Pretty())
303 + u.DOut("handleFindPeer: sending back '%s'\n", closest.ID.Pretty())
304 resp.Peers = []*peer.Peer{closest}
305 resp.Success = true
306 }
@@ -352,7 +357,7 @@ func (dht *IpfsDHT) Halt() {
357 }
358
359 func (dht *IpfsDHT) addProviderEntry(key u.Key, p *peer.Peer) {
355 - u.DOut("Adding %s as provider for '%s'", p.Key().Pretty(), key)
360 + u.DOut("Adding %s as provider for '%s'\n", p.Key().Pretty(), key)
361 dht.providerLock.Lock()
362 provs := dht.providers[key]
363 dht.providers[key] = append(provs, &providerInfo{time.Now(), p})
@@ -432,13 +437,13 @@ func (dht *IpfsDHT) getValueOrPeers(p *peer.Peer, key u.Key, timeout time.Durati
437 }
438 addr, err := ma.NewMultiaddr(pb.GetAddr())
439 if err != nil {
435 - u.PErr(err.Error())
440 + u.PErr("%v\n", err.Error())
441 continue
442 }
443
444 np, err := dht.network.GetConnection(peer.ID(pb.GetId()), addr)
445 if err != nil {
441 - u.PErr(err.Error())
446 + u.PErr("%v\n", err.Error())
447 continue
448 }
449
@@ -494,19 +499,19 @@ func (dht *IpfsDHT) getFromPeerList(key u.Key, timeout time.Duration,
499 if p == nil {
500 maddr, err := ma.NewMultiaddr(pinfo.GetAddr())
501 if err != nil {
497 - u.PErr("getValue error: %s", err)
502 + u.PErr("getValue error: %s\n", err)
503 continue
504 }
505
506 p, err = dht.network.GetConnection(peer.ID(pinfo.GetId()), maddr)
507 if err != nil {
503 - u.PErr("getValue error: %s", err)
508 + u.PErr("getValue error: %s\n", err)
509 continue
510 }
511 }
512 pmes, err := dht.getValueSingle(p, key, timeout, level)
513 if err != nil {
509 - u.DErr("getFromPeers error: %s", err)
514 + u.DErr("getFromPeers error: %s\n", err)
515 continue
516 }
517 dht.addProviderEntry(key, p)
@@ -520,6 +525,8 @@ func (dht *IpfsDHT) getFromPeerList(key u.Key, timeout time.Duration,
525 }
526
527 func (dht *IpfsDHT) getLocal(key u.Key) ([]byte, error) {
528 + dht.dslock.Lock()
529 + defer dht.dslock.Unlock()
530 v, err := dht.datastore.Get(ds.NewKey(string(key)))
531 if err != nil {
532 return nil, err
@@ -637,15 +644,15 @@ func (dht *IpfsDHT) addPeerList(key u.Key, peers []*PBDHTMessage_PBPeer) []*peer
644 // Dont add someone who is already on the list
645 p := dht.network.Find(u.Key(prov.GetId()))
646 if p == nil {
640 - u.DOut("given provider %s was not in our network already.", peer.ID(prov.GetId()).Pretty())
647 + u.DOut("given provider %s was not in our network already.\n", peer.ID(prov.GetId()).Pretty())
648 maddr, err := ma.NewMultiaddr(prov.GetAddr())
649 if err != nil {
643 - u.PErr("error connecting to new peer: %s", err)
650 + u.PErr("error connecting to new peer: %s\n", err)
651 continue
652 }
653 p, err = dht.network.GetConnection(peer.ID(prov.GetId()), maddr)
654 if err != nil {
648 - u.PErr("error connecting to new peer: %s", err)
655 + u.PErr("error connecting to new peer: %s\n", err)
656 continue
657 }
658 }
routing/dht/routing.go
+5 -6
@@ -30,8 +30,7 @@ var AlphaValue = 3
30 // GenerateMessageID creates and returns a new message ID
31 // TODO: determine a way of creating and managing message IDs
32 func GenerateMessageID() uint64 {
33 - //return (uint64(rand.Uint32()) << 32) & uint64(rand.Uint32())
34 - return uint64(rand.Uint32())
33 + return (uint64(rand.Uint32()) << 32) | uint64(rand.Uint32())
34 }
35
36 // This file implements the Routing interface for the IpfsDHT struct.
@@ -188,7 +187,7 @@ func (dht *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
187 }
188 val, peers, err := dht.getValueOrPeers(p, key, timeout/4, routeLevel)
189 if err != nil {
191 - u.DErr(err.Error())
190 + u.DErr("%v\n", err.Error())
191 c.Decrement()
192 continue
193 }
@@ -254,7 +253,7 @@ func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Pee
253 ll.EndLog()
254 ll.Print()
255 }()
257 - u.DOut("Find providers for: '%s'", key)
256 + u.DOut("Find providers for: '%s'\n", key)
257 p := dht.routingTables[0].NearestPeer(kb.ConvertKey(key))
258 if p == nil {
259 return nil, kb.ErrLookupFailure
@@ -288,7 +287,7 @@ func (dht *IpfsDHT) FindProviders(key u.Key, timeout time.Duration) ([]*peer.Pee
287
288 np, err := dht.network.GetConnection(peer.ID(closer[0].GetId()), maddr)
289 if err != nil {
291 - u.PErr("[%s] Failed to connect to: %s", dht.self.ID.Pretty(), closer[0].GetAddr())
290 + u.PErr("[%s] Failed to connect to: %s\n", dht.self.ID.Pretty(), closer[0].GetAddr())
291 level++
292 continue
293 }
@@ -361,7 +360,7 @@ func (dht *IpfsDHT) Ping(p *peer.Peer, timeout time.Duration) error {
360 case <-responseChan:
361 roundtrip := time.Since(before)
362 p.SetLatency(roundtrip)
364 - u.DOut("Ping took %s.", roundtrip.String())
363 + u.DOut("Ping took %s.\n", roundtrip.String())
364 return nil
365 case <-tout:
366 // Timed out, think about removing peer from network
swarm/conn.go
+2
@@ -40,6 +40,8 @@ func Dial(network string, peer *peer.Peer) (*Conn, error) {
40 return nil, err
41 }
42
43 + fmt.Printf("Making connection to: %s\n", host)
44 +
45 nconn, err := net.Dial(network, host)
46 if err != nil {
47 return nil, err
swarm/swarm.go
+11 -4
@@ -13,6 +13,8 @@ import (
13 ma "github.com/jbenet/go-multiaddr"
14 )
15
16 +var ErrAlreadyOpen = errors.New("Error: Connection to this peer already open.")
17 +
18 // Message represents a packet of information sent to or received from a
19 // particular Peer.
20 type Message struct {
@@ -27,7 +29,7 @@ type Message struct {
29 func NewMessage(p *peer.Peer, data proto.Message) *Message {
30 bytes, err := proto.Marshal(data)
31 if err != nil {
30 - u.PErr(err.Error())
32 + u.PErr("%v\n", err.Error())
33 return nil
34 }
35 return &Message{
@@ -162,7 +164,7 @@ func (s *Swarm) handleNewConn(nconn net.Conn) {
164
165 err := ident.Handshake(s.local, p, conn.Incoming.MsgChan, conn.Outgoing.MsgChan)
166 if err != nil {
165 - u.PErr(err.Error())
167 + u.PErr("%v\n", err.Error())
168 conn.Close()
169 return
170 }
@@ -228,9 +230,13 @@ func (s *Swarm) StartConn(conn *Conn) error {
230 return errors.New("Tried to start nil connection.")
231 }
232
231 - u.DOut("Starting connection: %s", conn.Peer.Key().Pretty())
233 + u.DOut("Starting connection: %s\n", conn.Peer.Key().Pretty())
234 // add to conns
235 s.connsLock.Lock()
236 + if _, ok := s.conns[conn.Peer.Key()]; ok {
237 + s.connsLock.Unlock()
238 + return ErrAlreadyOpen
239 + }
240 s.conns[conn.Peer.Key()] = conn
241 s.connsLock.Unlock()
242
@@ -249,7 +255,6 @@ func (s *Swarm) fanOut() {
255 if !ok {
256 return
257 }
252 - //u.DOut("fanOut: outgoing message for: '%s'", msg.Peer.Key().Pretty())
258
259 s.connsLock.RLock()
260 conn, found := s.conns[msg.Peer.Key()]
@@ -313,6 +318,8 @@ out:
318 }
319
320 func (s *Swarm) Find(key u.Key) *peer.Peer {
321 + s.connsLock.RLock()
322 + defer s.connsLock.RUnlock()
323 conn, found := s.conns[key]
324 if !found {
325 return nil