fixed get/put
Juan Batiz-Benet committed
Sep 19, 2014 at 08:07 UTC
043c09e14b6f654453e21a524933196cd454a367
5 files changed
+28
-18
routing/dht/dht.go
+10
-4
@@ -1,6 +1,7 @@
1
package dht
2
3
import (
4
+ "bytes"
5
"crypto/rand"
6
"errors"
7
"fmt"
@@ -190,15 +191,20 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p *peer.Peer, pmes *Message
191
return rpmes, nil
192
}
193
193
-func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p *peer.Peer, key string, value []byte) error {
194
+func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p *peer.Peer,
195
+ key string, value []byte) error {
196
+
197
pmes := newMessage(Message_PUT_VALUE, string(key), 0)
198
pmes.Value = value
196
-
197
- mes, err := msg.FromObject(p, pmes)
199
+ rpmes, err := dht.sendRequest(ctx, p, pmes)
200
if err != nil {
201
return err
202
}
201
- return dht.sender.SendMessage(ctx, mes)
203
+
204
+ if !bytes.Equal(rpmes.Value, pmes.Value) {
205
+ return errors.New("value not put correctly")
206
+ }
207
+ return nil
208
}
209
210
func (dht *IpfsDHT) putProvider(ctx context.Context, p *peer.Peer, key string) error {
routing/dht/dht_test.go
+1
-1
@@ -117,7 +117,7 @@ func TestPing(t *testing.T) {
117
}
118
119
func TestValueGetSet(t *testing.T) {
120
- u.Debug = false
120
+ u.Debug = true
121
addrA, err := ma.NewMultiaddr("/ip4/127.0.0.1/tcp/1235")
122
if err != nil {
123
t.Fatal(err)
routing/dht/handlers.go
+7
-6
@@ -38,7 +38,7 @@ func (dht *IpfsDHT) handlerForMsgType(t Message_MessageType) dhtHandler {
38
}
39
40
func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error) {
41
- u.DOut("handleGetValue for key: %s\n", pmes.GetKey())
41
+ u.DOut("[%s] handleGetValue for key: %s\n", dht.self.ID.Pretty(), pmes.GetKey())
42
43
// setup response
44
resp := newMessage(pmes.GetType(), pmes.GetKey(), pmes.GetClusterLevel())
@@ -50,11 +50,13 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
50
}
51
52
// let's first check if we have the value locally.
53
+ u.DOut("[%s] handleGetValue looking into ds\n", dht.self.ID.Pretty())
54
dskey := ds.NewKey(pmes.GetKey())
55
iVal, err := dht.datastore.Get(dskey)
56
+ u.DOut("[%s] handleGetValue looking into ds GOT %v\n", dht.self.ID.Pretty(), iVal)
57
58
// if we got an unexpected error, bail.
57
- if err != ds.ErrNotFound {
59
+ if err != nil && err != ds.ErrNotFound {
60
return nil, err
61
}
62
@@ -63,7 +65,7 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
65
66
// if we have the value, send it back
67
if err == nil {
66
- u.DOut("handleGetValue success!\n")
68
+ u.DOut("[%s] handleGetValue success!\n", dht.self.ID.Pretty())
69
70
byts, ok := iVal.([]byte)
71
if !ok {
@@ -85,7 +87,6 @@ func (dht *IpfsDHT) handleGetValue(p *peer.Peer, pmes *Message) (*Message, error
87
if closer != nil {
88
u.DOut("handleGetValue returning a closer peer: '%s'\n", closer.ID.Pretty())
89
resp.CloserPeers = peersToPBPeers([]*peer.Peer{closer})
88
- return resp, nil
90
}
91
92
return resp, nil
@@ -97,8 +98,8 @@ func (dht *IpfsDHT) handlePutValue(p *peer.Peer, pmes *Message) (*Message, error
98
defer dht.dslock.Unlock()
99
dskey := ds.NewKey(pmes.GetKey())
100
err := dht.datastore.Put(dskey, pmes.GetValue())
100
- u.DOut("[%s] handlePutValue %v %v", dht.self.ID.Pretty(), dskey, pmes.GetValue())
101
- return nil, err
101
+ u.DOut("[%s] handlePutValue %v %v\n", dht.self.ID.Pretty(), dskey, pmes.GetValue())
102
+ return pmes, err
103
}
104
105
func (dht *IpfsDHT) handlePing(p *peer.Peer, pmes *Message) (*Message, error) {
routing/dht/query.go
+8
-7
@@ -117,28 +117,29 @@ func (r *dhtQueryRunner) Run(peers []*peer.Peer) (*dhtQueryResult, error) {
117
// so workers are working.
118
119
// wait until they're done.
120
+ err := u.ErrNotFound
121
+
122
select {
123
case <-r.peersRemaining.Done():
124
r.cancel() // ran all and nothing. cancel all outstanding workers.
123
-
125
r.RLock()
126
defer r.RUnlock()
127
128
if len(r.errs) > 0 {
128
- return nil, r.errs[0]
129
+ err = r.errs[0]
130
}
130
- return nil, u.ErrNotFound
131
132
case <-r.ctx.Done():
133
r.RLock()
134
defer r.RUnlock()
135
+ err = r.ctx.Err()
136
+ }
137
136
- if r.result != nil && r.result.success {
137
- return r.result, nil
138
- }
139
- return nil, r.ctx.Err()
138
+ if r.result != nil && r.result.success {
139
+ return r.result, nil
140
}
141
142
+ return nil, err
143
}
144
145
func (r *dhtQueryRunner) addPeerToQuery(next *peer.Peer, benchmark *peer.Peer) {
routing/dht/routing.go
+2
@@ -89,6 +89,8 @@ func (dht *IpfsDHT) GetValue(key u.Key, timeout time.Duration) ([]byte, error) {
89
return nil, err
90
}
91
92
+ u.DOut("[%s] GetValue %v %v\n", dht.self.ID.Pretty(), key, result.value)
93
+
94
if result.value == nil {
95
return nil, u.ErrNotFound
96
}