change Provide RPC to not wait for an ACK, improves performance of 'Add' operations
Jeromy committed
Dec 16, 2014 at 04:01 UTC
2240272d850a245c55431215c84107fd69be921e
4 files changed
+30
-6
merkledag/merkledag.go
+3
-1
@@ -2,6 +2,7 @@
2
package merkledag
3
4
import (
5
+ "bytes"
6
"fmt"
7
"sync"
8
"time"
@@ -294,8 +295,9 @@ func FetchGraph(ctx context.Context, root *Node, serv DAGService) chan struct{}
295
// returns the indexes of any links pointing to it
296
func FindLinks(n *Node, k u.Key) []int {
297
var out []int
298
+ keybytes := []byte(k)
299
for i, lnk := range n.Links {
298
- if u.Key(lnk.Hash) == k {
300
+ if bytes.Equal([]byte(lnk.Hash), keybytes) {
301
out = append(out, i)
302
}
303
}
routing/dht/dht.go
+1
-4
@@ -120,15 +120,12 @@ func (dht *IpfsDHT) putProvider(ctx context.Context, p peer.Peer, key string) er
120
// add self as the provider
121
pmes.ProviderPeers = pb.PeersToPBPeers(dht.network, []peer.Peer{dht.self})
122
123
- rpmes, err := dht.sendRequest(ctx, p, pmes)
123
+ err := dht.sendMessage(ctx, p, pmes)
124
if err != nil {
125
return err
126
}
127
128
log.Debugf("%s putProvider: %s for %s", dht.self, p, u.Key(key))
129
- if rpmes.GetKey() != pmes.GetKey() {
130
- return errors.New("provider not added correctly")
131
- }
129
130
return nil
131
}
routing/dht/dht_net.go
+22
-1
@@ -8,8 +8,8 @@ import (
8
peer "github.com/jbenet/go-ipfs/peer"
9
pb "github.com/jbenet/go-ipfs/routing/dht/pb"
10
11
- ggio "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/io"
11
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
12
+ ggio "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/io"
13
)
14
15
// handleNewStream implements the inet.StreamHandler
@@ -102,3 +102,24 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *pb.Messa
102
log.Event(ctx, "dhtReceivedMessage", dht.self, p, rpmes)
103
return rpmes, nil
104
}
105
+
106
+// sendMessage sends out a message
107
+func (dht *IpfsDHT) sendMessage(ctx context.Context, p peer.Peer, pmes *pb.Message) error {
108
+
109
+ log.Debugf("%s dht starting stream", dht.self)
110
+ s, err := dht.network.NewStream(inet.ProtocolDHT, p)
111
+ if err != nil {
112
+ return err
113
+ }
114
+ defer s.Close()
115
+
116
+ w := ggio.NewDelimitedWriter(s)
117
+
118
+ log.Debugf("%s writing", dht.self)
119
+ if err := w.WriteMsg(pmes); err != nil {
120
+ return err
121
+ }
122
+ log.Event(ctx, "dhtSentMessage", dht.self, p, pmes)
123
+ log.Debugf("%s done", dht.self)
124
+ return nil
125
+}
routing/mock/mockrouting_test.go
+4
@@ -36,6 +36,9 @@ func TestClientFindProviders(t *testing.T) {
36
if err != nil {
37
t.Fatal(err)
38
}
39
+
40
+ // This is bad... but simulating networks is hard
41
+ time.Sleep(time.Millisecond * 300)
42
max := 100
43
44
providersFromHashTable, err := rs.Client(peer).FindProviders(context.Background(), k)
@@ -160,6 +163,7 @@ func TestValidAfter(t *testing.T) {
163
if err != nil {
164
t.Fatal(err)
165
}
166
+ t.Log("providers", providers)
167
if len(providers) != 1 {
168
t.Fail()
169
}