switch DHT entries over to be records, test currently fail
Jeromy committed
Nov 9, 2014 at 23:45 UTC
33985c530ed4bd8ccd4b93997cd67da58025ea20
9 files changed
+231
-31
namesys/publisher.go
+2
-2
@@ -49,7 +49,7 @@ func (p *ipnsPublisher) Publish(k ci.PrivKey, value string) error {
49
50
nameb := u.Hash(pkbytes)
51
namekey := u.Key(nameb).Pretty()
52
- ipnskey := u.Hash([]byte("/ipns/" + namekey))
52
+ ipnskey := []byte("/ipns/" + namekey)
53
54
// Store associated public key
55
timectx, _ := context.WithDeadline(ctx, time.Now().Add(time.Second*4))
@@ -58,7 +58,7 @@ func (p *ipnsPublisher) Publish(k ci.PrivKey, value string) error {
58
return err
59
}
60
61
- // Store ipns entry at h("/ipns/"+b58(h(pubkey)))
61
+ // Store ipns entry at "/ipns/"+b58(h(pubkey))
62
timectx, _ = context.WithDeadline(ctx, time.Now().Add(time.Second*4))
63
err = p.routing.PutValue(timectx, u.Key(ipnskey), data)
64
if err != nil {
namesys/routing.go
+1
-1
@@ -46,7 +46,7 @@ func (r *routingResolver) Resolve(name string) (string, error) {
46
47
// use the routing system to get the name.
48
// /ipns/<name>
49
- h := u.Hash([]byte("/ipns/" + name))
49
+ h := []byte("/ipns/" + name)
50
51
ipnsKey := u.Key(h)
52
val, err := r.routing.GetValue(ctx, ipnsKey)
routing/dht/dht.go
+47
-10
@@ -60,6 +60,9 @@ type IpfsDHT struct {
60
//lock to make diagnostics work better
61
diaglock sync.Mutex
62
63
+ // record validator funcs
64
+ Validators map[string]ValidatorFunc
65
+
66
ctxc.ContextCloser
67
}
68
@@ -81,6 +84,7 @@ func NewDHT(ctx context.Context, p peer.Peer, ps peer.Peerstore, dialer inet.Dia
84
dht.routingTables[1] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Millisecond*1000)
85
dht.routingTables[2] = kb.NewRoutingTable(20, kb.ConvertPeerID(p.ID()), time.Hour)
86
dht.birth = time.Now()
87
+ dht.Validators = make(map[string]ValidatorFunc)
88
89
if doPinging {
90
dht.Children().Add(1)
@@ -215,16 +219,16 @@ func (dht *IpfsDHT) sendRequest(ctx context.Context, p peer.Peer, pmes *pb.Messa
219
220
// putValueToNetwork stores the given key/value pair at the peer 'p'
221
func (dht *IpfsDHT) putValueToNetwork(ctx context.Context, p peer.Peer,
218
- key string, value []byte) error {
222
+ key string, rec *pb.Record) error {
223
224
pmes := pb.NewMessage(pb.Message_PUT_VALUE, string(key), 0)
221
- pmes.Value = value
225
+ pmes.Record = rec
226
rpmes, err := dht.sendRequest(ctx, p, pmes)
227
if err != nil {
228
return err
229
}
230
227
- if !bytes.Equal(rpmes.Value, pmes.Value) {
231
+ if !bytes.Equal(rpmes.GetRecord().Value, pmes.GetRecord().Value) {
232
return errors.New("value not put correctly")
233
}
234
return nil
@@ -260,11 +264,16 @@ func (dht *IpfsDHT) getValueOrPeers(ctx context.Context, p peer.Peer,
264
return nil, nil, err
265
}
266
263
- log.Debugf("pmes.GetValue() %v", pmes.GetValue())
264
- if value := pmes.GetValue(); value != nil {
267
+ if record := pmes.GetRecord(); record != nil {
268
// Success! We were given the value
269
log.Debug("getValueOrPeers: got value")
267
- return value, nil, nil
270
+
271
+ // make sure record is still valid
272
+ err = dht.verifyRecord(record)
273
+ if err != nil {
274
+ return nil, nil, err
275
+ }
276
+ return record.GetValue(), nil, nil
277
}
278
279
// TODO decide on providers. This probably shouldn't be happening.
@@ -325,10 +334,15 @@ func (dht *IpfsDHT) getFromPeerList(ctx context.Context, key u.Key,
334
continue
335
}
336
328
- if value := pmes.GetValue(); value != nil {
337
+ if record := pmes.GetRecord(); record != nil {
338
// Success! We were given the value
339
+
340
+ err := dht.verifyRecord(record)
341
+ if err != nil {
342
+ return nil, err
343
+ }
344
dht.providers.AddProvider(key, p)
331
- return value, nil
345
+ return record.GetValue(), nil
346
}
347
}
348
return nil, routing.ErrNotFound
@@ -347,12 +361,35 @@ func (dht *IpfsDHT) getLocal(key u.Key) ([]byte, error) {
361
if !ok {
362
return nil, errors.New("value stored in datastore not []byte")
363
}
350
- return byt, nil
364
+ rec := new(pb.Record)
365
+ err = proto.Unmarshal(byt, rec)
366
+ if err != nil {
367
+ return nil, err
368
+ }
369
+
370
+ // TODO: 'if paranoid'
371
+ if u.Debug {
372
+ err = dht.verifyRecord(rec)
373
+ if err != nil {
374
+ return nil, err
375
+ }
376
+ }
377
+
378
+ return rec.GetValue(), nil
379
}
380
381
// putLocal stores the key value pair in the datastore
382
func (dht *IpfsDHT) putLocal(key u.Key, value []byte) error {
355
- return dht.datastore.Put(key.DsKey(), value)
383
+ rec, err := dht.makePutRecord(key, value)
384
+ if err != nil {
385
+ return err
386
+ }
387
+ data, err := proto.Marshal(rec)
388
+ if err != nil {
389
+ return err
390
+ }
391
+
392
+ return dht.datastore.Put(key.DsKey(), data)
393
}
394
395
// Update signals to all routingTables to Update their last-seen status
routing/dht/ext_test.go
+13
-9
@@ -124,10 +124,10 @@ func TestGetFailures(t *testing.T) {
124
fs := &fauxSender{}
125
126
peerstore := peer.NewPeerstore()
127
- local := peer.WithIDString("test_peer")
127
+ local := makePeer(nil)
128
129
d := NewDHT(ctx, local, peerstore, fn, fs, ds.NewMapDatastore())
130
- other := peer.WithIDString("other_peer")
130
+ other := makePeer(nil)
131
d.Update(other)
132
133
// This one should time out
@@ -173,10 +173,14 @@ func TestGetFailures(t *testing.T) {
173
// Now we test this DHT's handleGetValue failure
174
typ := pb.Message_GET_VALUE
175
str := "hello"
176
+ rec, err := d.makePutRecord(u.Key(str), []byte("blah"))
177
+ if err != nil {
178
+ t.Fatal(err)
179
+ }
180
req := pb.Message{
177
- Type: &typ,
178
- Key: &str,
179
- Value: []byte{0},
181
+ Type: &typ,
182
+ Key: &str,
183
+ Record: rec,
184
}
185
186
// u.POut("handleGetValue Test\n")
@@ -192,10 +196,10 @@ func TestGetFailures(t *testing.T) {
196
if err != nil {
197
t.Fatal(err)
198
}
195
- if pmes.GetValue() != nil {
199
+ if pmes.GetRecord() != nil {
200
t.Fatal("shouldnt have value")
201
}
198
- if pmes.GetCloserPeers() != nil {
202
+ if len(pmes.GetCloserPeers()) > 0 {
203
t.Fatal("shouldnt have closer peers")
204
}
205
if pmes.GetProviderPeers() != nil {
@@ -221,7 +225,7 @@ func TestNotFound(t *testing.T) {
225
fn := &fauxNet{}
226
fs := &fauxSender{}
227
224
- local := peer.WithIDString("test_peer")
228
+ local := makePeer(nil)
229
peerstore := peer.NewPeerstore()
230
peerstore.Add(local)
231
@@ -287,7 +291,7 @@ func TestLessThanKResponses(t *testing.T) {
291
u.Debug = false
292
fn := &fauxNet{}
293
fs := &fauxSender{}
290
- local := peer.WithIDString("test_peer")
294
+ local := makePeer(nil)
295
peerstore := peer.NewPeerstore()
296
peerstore.Add(local)
297
routing/dht/handlers.go
+24
-3
@@ -5,6 +5,8 @@ import (
5
"fmt"
6
"time"
7
8
+ "code.google.com/p/goprotobuf/proto"
9
+
10
peer "github.com/jbenet/go-ipfs/peer"
11
pb "github.com/jbenet/go-ipfs/routing/dht/pb"
12
u "github.com/jbenet/go-ipfs/util"
@@ -72,7 +74,14 @@ func (dht *IpfsDHT) handleGetValue(p peer.Peer, pmes *pb.Message) (*pb.Message,
74
return nil, fmt.Errorf("datastore had non byte-slice value for %v", dskey)
75
}
76
75
- resp.Value = byts
77
+ rec := new(pb.Record)
78
+ err := proto.Unmarshal(byts, rec)
79
+ if err != nil {
80
+ log.Error("Failed to unmarshal dht record from datastore")
81
+ return nil, err
82
+ }
83
+
84
+ resp.Record = rec
85
}
86
87
// if we know any providers for the requested value, return those.
@@ -102,8 +111,20 @@ func (dht *IpfsDHT) handlePutValue(p peer.Peer, pmes *pb.Message) (*pb.Message,
111
dht.dslock.Lock()
112
defer dht.dslock.Unlock()
113
dskey := u.Key(pmes.GetKey()).DsKey()
105
- err := dht.datastore.Put(dskey, pmes.GetValue())
106
- log.Debugf("%s handlePutValue %v %v\n", dht.self, dskey, pmes.GetValue())
114
+
115
+ err := dht.verifyRecord(pmes.GetRecord())
116
+ if err != nil {
117
+ log.Error("Bad dht record in put request")
118
+ return nil, err
119
+ }
120
+
121
+ data, err := proto.Marshal(pmes.GetRecord())
122
+ if err != nil {
123
+ return nil, err
124
+ }
125
+
126
+ err = dht.datastore.Put(dskey, data)
127
+ log.Debugf("%s handlePutValue %v\n", dht.self, dskey)
128
return pmes, err
129
}
130
routing/dht/pb/dht.pb.go
+51
-4
@@ -10,10 +10,11 @@ It is generated from these files:
10
11
It has these top-level messages:
12
Message
13
+ Record
14
*/
15
package dht_pb
16
16
-import proto "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/gogoprotobuf/proto"
17
+import proto "code.google.com/p/gogoprotobuf/proto"
18
import math "math"
19
20
// Reference imports to suppress errors if they are not otherwise used.
@@ -75,7 +76,7 @@ type Message struct {
76
Key *string `protobuf:"bytes,2,opt,name=key" json:"key,omitempty"`
77
// Used to return a value
78
// PUT_VALUE, GET_VALUE
78
- Value []byte `protobuf:"bytes,3,opt,name=value" json:"value,omitempty"`
79
+ Record *Record `protobuf:"bytes,3,opt,name=record" json:"record,omitempty"`
80
// Used to return peers closer to a key in a query
81
// GET_VALUE, GET_PROVIDERS, FIND_NODE
82
CloserPeers []*Message_Peer `protobuf:"bytes,8,rep,name=closerPeers" json:"closerPeers,omitempty"`
@@ -110,9 +111,9 @@ func (m *Message) GetKey() string {
111
return ""
112
}
113
113
-func (m *Message) GetValue() []byte {
114
+func (m *Message) GetRecord() *Record {
115
if m != nil {
115
- return m.Value
116
+ return m.Record
117
}
118
return nil
119
}
@@ -155,6 +156,52 @@ func (m *Message_Peer) GetAddr() string {
156
return ""
157
}
158
159
+// Record represents a dht record that contains a value
160
+// for a key value pair
161
+type Record struct {
162
+ // The key that references this record
163
+ Key *string `protobuf:"bytes,1,opt,name=key" json:"key,omitempty"`
164
+ // The actual value this record is storing
165
+ Value []byte `protobuf:"bytes,2,opt,name=value" json:"value,omitempty"`
166
+ // hash of the authors public key
167
+ Author *string `protobuf:"bytes,3,opt,name=author" json:"author,omitempty"`
168
+ // A PKI signature for the key+value+author
169
+ Signature []byte `protobuf:"bytes,4,opt,name=signature" json:"signature,omitempty"`
170
+ XXX_unrecognized []byte `json:"-"`
171
+}
172
+
173
+func (m *Record) Reset() { *m = Record{} }
174
+func (m *Record) String() string { return proto.CompactTextString(m) }
175
+func (*Record) ProtoMessage() {}
176
+
177
+func (m *Record) GetKey() string {
178
+ if m != nil && m.Key != nil {
179
+ return *m.Key
180
+ }
181
+ return ""
182
+}
183
+
184
+func (m *Record) GetValue() []byte {
185
+ if m != nil {
186
+ return m.Value
187
+ }
188
+ return nil
189
+}
190
+
191
+func (m *Record) GetAuthor() string {
192
+ if m != nil && m.Author != nil {
193
+ return *m.Author
194
+ }
195
+ return ""
196
+}
197
+
198
+func (m *Record) GetSignature() []byte {
199
+ if m != nil {
200
+ return m.Signature
201
+ }
202
+ return nil
203
+}
204
+
205
func init() {
206
proto.RegisterEnum("dht.pb.Message_MessageType", Message_MessageType_name, Message_MessageType_value)
207
}
routing/dht/pb/dht.proto
+17
-1
@@ -29,7 +29,7 @@ message Message {
29
30
// Used to return a value
31
// PUT_VALUE, GET_VALUE
32
- optional bytes value = 3;
32
+ optional Record record = 3;
33
34
// Used to return peers closer to a key in a query
35
// GET_VALUE, GET_PROVIDERS, FIND_NODE
@@ -39,3 +39,19 @@ message Message {
39
// GET_VALUE, ADD_PROVIDER, GET_PROVIDERS
40
repeated Peer providerPeers = 9;
41
}
42
+
43
+// Record represents a dht record that contains a value
44
+// for a key value pair
45
+message Record {
46
+ // The key that references this record
47
+ optional string key = 1;
48
+
49
+ // The actual value this record is storing
50
+ optional bytes value = 2;
51
+
52
+ // hash of the authors public key
53
+ optional string author = 3;
54
+
55
+ // A PKI signature for the key+value+author
56
+ optional bytes signature = 4;
57
+}
routing/dht/records.go
new
+69
@@ -0,0 +1,69 @@
1
+package dht
2
+
3
+import (
4
+ "bytes"
5
+ "errors"
6
+ "strings"
7
+
8
+ "code.google.com/p/goprotobuf/proto"
9
+ "github.com/jbenet/go-ipfs/peer"
10
+ pb "github.com/jbenet/go-ipfs/routing/dht/pb"
11
+ u "github.com/jbenet/go-ipfs/util"
12
+)
13
+
14
+type ValidatorFunc func(u.Key, []byte) error
15
+
16
+var ErrBadRecord = errors.New("bad dht record")
17
+var ErrInvalidRecordType = errors.New("invalid record keytype")
18
+
19
+// creates and signs a dht record for the given key/value pair
20
+func (dht *IpfsDHT) makePutRecord(key u.Key, value []byte) (*pb.Record, error) {
21
+ record := new(pb.Record)
22
+
23
+ record.Key = proto.String(key.String())
24
+ record.Value = value
25
+ record.Author = proto.String(string(dht.self.ID()))
26
+ blob := bytes.Join([][]byte{[]byte(key), value, []byte(dht.self.ID())}, []byte{})
27
+ sig, err := dht.self.PrivKey().Sign(blob)
28
+ if err != nil {
29
+ return nil, err
30
+ }
31
+ record.Signature = sig
32
+ return record, nil
33
+}
34
+
35
+func (dht *IpfsDHT) verifyRecord(r *pb.Record) error {
36
+ // First, validate the signature
37
+ p, err := dht.peerstore.Get(peer.ID(r.GetAuthor()))
38
+ if err != nil {
39
+ return err
40
+ }
41
+
42
+ blob := bytes.Join([][]byte{[]byte(r.GetKey()),
43
+ r.GetValue(),
44
+ []byte(r.GetKey())}, []byte{})
45
+
46
+ ok, err := p.PubKey().Verify(blob, r.GetSignature())
47
+ if err != nil {
48
+ return err
49
+ }
50
+
51
+ if !ok {
52
+ return ErrBadRecord
53
+ }
54
+
55
+ // Now, check validity func
56
+ parts := strings.Split(r.GetKey(), "/")
57
+ if len(parts) < 2 {
58
+ log.Error("Record had bad key: %s", r.GetKey())
59
+ return ErrBadRecord
60
+ }
61
+
62
+ fnc, ok := dht.Validators[parts[0]]
63
+ if !ok {
64
+ log.Errorf("Unrecognized key prefix: %s", parts[0])
65
+ return ErrInvalidRecordType
66
+ }
67
+
68
+ return fnc(u.Key(r.GetKey()), r.GetValue())
69
+}
routing/dht/routing.go
+7
-1
@@ -25,6 +25,12 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
25
return err
26
}
27
28
+ rec, err := dht.makePutRecord(key, value)
29
+ if err != nil {
30
+ log.Error("Creation of record failed!")
31
+ return err
32
+ }
33
+
34
var peers []peer.Peer
35
for _, route := range dht.routingTables {
36
npeers := route.NearestPeers(kb.ConvertKey(key), KValue)
@@ -33,7 +39,7 @@ func (dht *IpfsDHT) PutValue(ctx context.Context, key u.Key, value []byte) error
39
40
query := newQuery(key, dht.dialer, func(ctx context.Context, p peer.Peer) (*dhtQueryResult, error) {
41
log.Debugf("%s PutValue qry part %v", dht.self, p)
36
- err := dht.putValueToNetwork(ctx, p, string(key), value)
42
+ err := dht.putValueToNetwork(ctx, p, string(key), rec)
43
if err != nil {
44
return nil, err
45
}