@cryptotaxi247 / kubo / commits / 4b01ad0b7

fixed pin json marshal

Juan Batiz-Benet committed Oct 22, 2014 at 02:55 UTC 4b01ad0b7c7f65b8b4df700aef28ed4a1f55d247
4 files changed +62 -60
blocks/set/set.go
+2 -19
@@ -1,10 +1,6 @@
1 package set
2
3 import (
4 - "errors"
5 -
6 - ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
7 -
4 "github.com/jbenet/go-ipfs/blocks/bloom"
5 "github.com/jbenet/go-ipfs/util"
6 )
@@ -28,19 +24,6 @@ func SimpleSetFromKeys(keys []util.Key) BlockSet {
24 return sbs
25 }
26
31 -func SetFromDatastore(d ds.Datastore, k ds.Key) (BlockSet, error) {
32 - ikeys, err := d.Get(k)
33 - if err != nil {
34 - return nil, err
35 - }
36 -
37 - keys, ok := ikeys.([]util.Key)
38 - if !ok {
39 - return nil, errors.New("Incorrect type for keys from datastore")
40 - }
41 - return SimpleSetFromKeys(keys), nil
42 -}
43 -
27 func NewSimpleBlockSet() BlockSet {
28 return &simpleBlockSet{blocks: make(map[util.Key]struct{})}
29 }
@@ -64,7 +47,7 @@ func (b *simpleBlockSet) HasKey(k util.Key) bool {
47
48 func (b *simpleBlockSet) GetBloomFilter() bloom.Filter {
49 f := bloom.BasicFilter()
67 - for k, _ := range b.blocks {
50 + for k := range b.blocks {
51 f.Add([]byte(k))
52 }
53 return f
@@ -72,7 +55,7 @@ func (b *simpleBlockSet) GetBloomFilter() bloom.Filter {
55
56 func (b *simpleBlockSet) GetKeys() []util.Key {
57 var out []util.Key
75 - for k, _ := range b.blocks {
58 + for k := range b.blocks {
59 out = append(out, k)
60 }
61 return out
pin/indirect.go
+16 -8
@@ -1,8 +1,6 @@
1 package pin
2
3 import (
4 - "errors"
5 -
4 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
5 "github.com/jbenet/go-ipfs/blocks/set"
6 "github.com/jbenet/go-ipfs/util"
@@ -21,23 +19,33 @@ func NewIndirectPin(dstore ds.Datastore) *indirectPin {
19 }
20
21 func loadIndirPin(d ds.Datastore, k ds.Key) (*indirectPin, error) {
24 - irefcnt, err := d.Get(k)
22 + var rcStore map[string]int
23 + err := loadSet(d, k, &rcStore)
24 if err != nil {
25 return nil, err
26 }
28 - refcnt, ok := irefcnt.(map[util.Key]int)
29 - if !ok {
30 - return nil, errors.New("invalid type from datastore")
31 - }
27
28 + refcnt := make(map[util.Key]int)
29 var keys []util.Key
34 - for k, _ := range refcnt {
30 + for encK, v := range rcStore {
31 + k := util.B58KeyDecode(encK)
32 keys = append(keys, k)
33 + refcnt[k] = v
34 }
35 + log.Debug("indirPin keys: %#v", keys)
36
37 return &indirectPin{blockset: set.SimpleSetFromKeys(keys), refCounts: refcnt}, nil
38 }
39
40 +func storeIndirPin(d ds.Datastore, k ds.Key, p *indirectPin) error {
41 +
42 + rcStore := map[string]int{}
43 + for k, v := range p.refCounts {
44 + rcStore[util.B58KeyEncode(k)] = v
45 + }
46 + return storeSet(d, k, rcStore)
47 +}
48 +
49 func (i *indirectPin) Increment(k util.Key) {
50 c := i.refCounts[k]
51 i.refCounts[k] = c + 1
pin/pin.go
+41 -30
@@ -3,8 +3,9 @@ package pin
3 import (
4
5 //ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
6 - "bytes"
6 +
7 "encoding/json"
8 + "errors"
9 "sync"
10
11 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
@@ -14,6 +15,7 @@ import (
15 "github.com/jbenet/go-ipfs/util"
16 )
17
18 +var log = util.Logger("pin")
19 var recursePinDatastoreKey = ds.NewKey("/local/pins/recursive/keys")
20 var directPinDatastoreKey = ds.NewKey("/local/pins/direct/keys")
21 var indirectPinDatastoreKey = ds.NewKey("/local/pins/indirect/keys")
@@ -89,9 +91,8 @@ func (p *pinner) Unpin(k util.Key, recurse bool) error {
91 }
92
93 return p.unpinLinks(node)
92 - } else {
93 - p.directPin.RemoveBlock(k)
94 }
95 + p.directPin.RemoveBlock(k)
96 return nil
97 }
98
@@ -153,21 +154,31 @@ func (p *pinner) IsPinned(key util.Key) bool {
154 func LoadPinner(d ds.Datastore, dserv *mdag.DAGService) (Pinner, error) {
155 p := new(pinner)
156
156 - var err error
157 - p.recursePin, err = set.SetFromDatastore(d, recursePinDatastoreKey)
158 - if err != nil {
159 - return nil, err
157 + { // load recursive set
158 + var recurseKeys []util.Key
159 + if err := loadSet(d, recursePinDatastoreKey, &recurseKeys); err != nil {
160 + return nil, err
161 + }
162 + p.recursePin = set.SimpleSetFromKeys(recurseKeys)
163 }
161 - p.directPin, err = set.SetFromDatastore(d, directPinDatastoreKey)
162 - if err != nil {
163 - return nil, err
164 +
165 + { // load direct set
166 + var directKeys []util.Key
167 + if err := loadSet(d, directPinDatastoreKey, &directKeys); err != nil {
168 + return nil, err
169 + }
170 + p.directPin = set.SimpleSetFromKeys(directKeys)
171 }
172
166 - p.indirPin, err = loadIndirPin(d, indirectPinDatastoreKey)
167 - if err != nil {
168 - return nil, err
173 + { // load indirect set
174 + var err error
175 + p.indirPin, err = loadIndirPin(d, indirectPinDatastoreKey)
176 + if err != nil {
177 + return nil, err
178 + }
179 }
180
181 + // assign services
182 p.dserv = dserv
183 p.dstore = d
184
@@ -177,43 +188,43 @@ func LoadPinner(d ds.Datastore, dserv *mdag.DAGService) (Pinner, error) {
188 func (p *pinner) Flush() error {
189 p.lock.RLock()
190 defer p.lock.RUnlock()
180 - buf := new(bytes.Buffer)
181 - enc := json.NewEncoder(buf)
191
183 - recurse := p.recursePin.GetKeys()
184 - err := enc.Encode(recurse)
192 + err := storeSet(p.dstore, directPinDatastoreKey, p.directPin.GetKeys())
193 if err != nil {
194 return err
195 }
196
189 - err = p.dstore.Put(recursePinDatastoreKey, buf.Bytes())
197 + err = storeSet(p.dstore, recursePinDatastoreKey, p.recursePin.GetKeys())
198 if err != nil {
199 return err
200 }
201
194 - buf = new(bytes.Buffer)
195 - enc = json.NewEncoder(buf)
196 - direct := p.directPin.GetKeys()
197 - err = enc.Encode(direct)
202 + err = storeIndirPin(p.dstore, indirectPinDatastoreKey, p.indirPin)
203 if err != nil {
204 return err
205 }
206 + return nil
207 +}
208
202 - err = p.dstore.Put(directPinDatastoreKey, buf.Bytes())
209 +// helpers to marshal / unmarshal a pin set
210 +func storeSet(d ds.Datastore, k ds.Key, val interface{}) error {
211 + buf, err := json.Marshal(val)
212 if err != nil {
213 return err
214 }
215
207 - buf = new(bytes.Buffer)
208 - enc = json.NewEncoder(buf)
209 - err = enc.Encode(p.indirPin.refCounts)
216 + return d.Put(k, buf)
217 +}
218 +
219 +func loadSet(d ds.Datastore, k ds.Key, val interface{}) error {
220 + buf, err := d.Get(k)
221 if err != nil {
222 return err
223 }
224
214 - err = p.dstore.Put(indirectPinDatastoreKey, buf.Bytes())
215 - if err != nil {
216 - return err
225 + bf, ok := buf.([]byte)
226 + if !ok {
227 + return errors.New("invalid pin set value in datastore")
228 }
218 - return nil
229 + return json.Unmarshal(bf, val)
230 }
pin/pin_test.go
+3 -3
@@ -3,7 +3,7 @@ package pin
3 import (
4 "testing"
5
6 - "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
6 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
7 bs "github.com/jbenet/go-ipfs/blockservice"
8 mdag "github.com/jbenet/go-ipfs/merkledag"
9 "github.com/jbenet/go-ipfs/util"
@@ -18,7 +18,7 @@ func randNode() (*mdag.Node, util.Key) {
18 }
19
20 func TestPinnerBasic(t *testing.T) {
21 - dstore := datastore.NewMapDatastore()
21 + dstore := ds.NewMapDatastore()
22 bserv, err := bs.NewBlockService(dstore, nil)
23 if err != nil {
24 t.Fatal(err)
@@ -103,7 +103,7 @@ func TestPinnerBasic(t *testing.T) {
103
104 // c should still be pinned under b
105 if !p.IsPinned(ck) {
106 - t.Fatal("Recursive unpin fail.")
106 + t.Fatal("Recursive / indirect unpin fail.")
107 }
108
109 err = p.Flush()