@cryptotaxi247 / kubo / commits / 6458ddcd3

flesh out pinning object, needs tests and cli wiring still

Jeromy committed Oct 19, 2014 at 15:28 UTC 6458ddcd32291b4e1fef79c7b5d2e0c080fc0cfe
13 files changed +370 -70
Godeps/Godeps.json
+1 -1
@@ -1,6 +1,6 @@
1 {
2 "ImportPath": "github.com/jbenet/go-ipfs",
3 - "GoVersion": "go1.3",
3 + "GoVersion": "go1.3.3",
4 "Packages": [
5 "./..."
6 ],
Godeps/_workspace/src/github.com/jbenet/datastore.go/keytransform/doc.go new
+25
@@ -0,0 +1,25 @@
1 +// Package keytransform introduces a Datastore Shim that transforms keys before
2 +// passing them to its child. It can be used to manipulate what keys look like
3 +// to the user, for example namespacing keys, reversing them, etc.
4 +//
5 +// Use the Wrap function to wrap a datastore with any KeyTransform.
6 +// A KeyTransform is simply an interface with two functions, a conversion and
7 +// its inverse. For example:
8 +//
9 +// import (
10 +// ktds "github.com/jbenet/datastore.go/keytransform"
11 +// ds "github.com/jbenet/datastore.go"
12 +// )
13 +//
14 +// func reverseKey(k ds.Key) ds.Key {
15 +// return k.Reverse()
16 +// }
17 +//
18 +// func invertKeys(d ds.Datastore) {
19 +// return ktds.Wrap(d, &ktds.Pair{
20 +// Convert: reverseKey,
21 +// Invert: reverseKey, // reverse is its own inverse.
22 +// })
23 +// }
24 +//
25 +package keytransform
Godeps/_workspace/src/github.com/jbenet/datastore.go/keytransform/interface.go new
+34
@@ -0,0 +1,34 @@
1 +package keytransform
2 +
3 +import ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
4 +
5 +// KeyMapping is a function that maps one key to annother
6 +type KeyMapping func(ds.Key) ds.Key
7 +
8 +// KeyTransform is an object with a pair of functions for (invertibly)
9 +// transforming keys
10 +type KeyTransform interface {
11 + ConvertKey(ds.Key) ds.Key
12 + InvertKey(ds.Key) ds.Key
13 +}
14 +
15 +// Datastore is a keytransform.Datastore
16 +type Datastore interface {
17 + ds.Shim
18 + KeyTransform
19 +}
20 +
21 +// Wrap wraps a given datastore with a KeyTransform function.
22 +// The resulting wrapped datastore will use the transform on all Datastore
23 +// operations.
24 +func Wrap(child ds.Datastore, t KeyTransform) Datastore {
25 + if t == nil {
26 + panic("t (KeyTransform) is nil")
27 + }
28 +
29 + if child == nil {
30 + panic("child (ds.Datastore) is nil")
31 + }
32 +
33 + return &ktds{child: child, KeyTransform: t}
34 +}
Godeps/_workspace/src/github.com/jbenet/datastore.go/namespace/doc.go new
+24
@@ -0,0 +1,24 @@
1 +// Package namespace introduces a namespace Datastore Shim, which basically
2 +// mounts the entire child datastore under a prefix.
3 +//
4 +// Use the Wrap function to wrap a datastore with any Key prefix. For example:
5 +//
6 +// import (
7 +// "fmt"
8 +//
9 +// ds "github.com/jbenet/datastore.go"
10 +// nsds "github.com/jbenet/datastore.go/namespace"
11 +// )
12 +//
13 +// func main() {
14 +// mp := ds.NewMapDatastore()
15 +// ns := nsds.Wrap(mp, ds.NewKey("/foo/bar"))
16 +//
17 +// // in the Namespace Datastore:
18 +// ns.Put(ds.NewKey("/beep"), "boop")
19 +// v2, _ := ns.Get(ds.NewKey("/beep")) // v2 == "boop"
20 +//
21 +// // and, in the underlying MapDatastore:
22 +// v3, _ := mp.Get(ds.NewKey("/foo/bar/beep")) // v3 == "boop"
23 +// }
24 +package namespace
Godeps/_workspace/src/github.com/jbenet/datastore.go/namespace/example_test.go new
+30
@@ -0,0 +1,30 @@
1 +package namespace_test
2 +
3 +import (
4 + "fmt"
5 +
6 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
7 + nsds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go/namespace"
8 +)
9 +
10 +func Example() {
11 + mp := ds.NewMapDatastore()
12 + ns := nsds.Wrap(mp, ds.NewKey("/foo/bar"))
13 +
14 + k := ds.NewKey("/beep")
15 + v := "boop"
16 +
17 + ns.Put(k, v)
18 + fmt.Printf("ns.Put %s %s\n", k, v)
19 +
20 + v2, _ := ns.Get(k)
21 + fmt.Printf("ns.Get %s -> %s\n", k, v2)
22 +
23 + k3 := ds.NewKey("/foo/bar/beep")
24 + v3, _ := mp.Get(k3)
25 + fmt.Printf("mp.Get %s -> %s\n", k3, v3)
26 + // Output:
27 + // ns.Put /beep -> boop
28 + // ns.Get /beep -> boop
29 + // mp.Get /foo/bar/beep -> boop
30 +}
Godeps/_workspace/src/github.com/jbenet/datastore.go/namespace/namespace.go new
+44
@@ -0,0 +1,44 @@
1 +package namespace
2 +
3 +import (
4 + "fmt"
5 + "strings"
6 +
7 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
8 + ktds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go/keytransform"
9 +)
10 +
11 +// PrefixTransform constructs a KeyTransform with a pair of functions that
12 +// add or remove the given prefix key.
13 +//
14 +// Warning: will panic if prefix not found when it should be there. This is
15 +// to avoid insidious data inconsistency errors.
16 +func PrefixTransform(prefix ds.Key) ktds.KeyTransform {
17 + return &ktds.Pair{
18 +
19 + // Convert adds the prefix
20 + Convert: func(k ds.Key) ds.Key {
21 + return prefix.Child(k)
22 + },
23 +
24 + // Invert removes the prefix. panics if prefix not found.
25 + Invert: func(k ds.Key) ds.Key {
26 + if !prefix.IsAncestorOf(k) {
27 + fmt.Errorf("Expected prefix (%s) in key (%s)", prefix, k)
28 + panic("expected prefix not found")
29 + }
30 +
31 + s := strings.TrimPrefix(k.String(), prefix.String())
32 + return ds.NewKey(s)
33 + },
34 + }
35 +}
36 +
37 +// Wrap wraps a given datastore with a key-prefix.
38 +func Wrap(child ds.Datastore, prefix ds.Key) ktds.Datastore {
39 + if child == nil {
40 + panic("child (ds.Datastore) is nil")
41 + }
42 +
43 + return ktds.Wrap(child, PrefixTransform(prefix))
44 +}
Godeps/_workspace/src/github.com/jbenet/datastore.go/namespace/namespace_test.go new
+72
@@ -0,0 +1,72 @@
1 +package namespace_test
2 +
3 +import (
4 + "bytes"
5 + "sort"
6 + "testing"
7 +
8 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
9 + ns "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go/namespace"
10 + . "launchpad.net/gocheck"
11 +)
12 +
13 +// Hook up gocheck into the "go test" runner.
14 +func Test(t *testing.T) { TestingT(t) }
15 +
16 +type DSSuite struct{}
17 +
18 +var _ = Suite(&DSSuite{})
19 +
20 +func (ks *DSSuite) TestBasic(c *C) {
21 +
22 + mpds := ds.NewMapDatastore()
23 + nsds := ns.Wrap(mpds, ds.NewKey("abc"))
24 +
25 + keys := strsToKeys([]string{
26 + "foo",
27 + "foo/bar",
28 + "foo/bar/baz",
29 + "foo/barb",
30 + "foo/bar/bazb",
31 + "foo/bar/baz/barb",
32 + })
33 +
34 + for _, k := range keys {
35 + err := nsds.Put(k, []byte(k.String()))
36 + c.Check(err, Equals, nil)
37 + }
38 +
39 + for _, k := range keys {
40 + v1, err := nsds.Get(k)
41 + c.Check(err, Equals, nil)
42 + c.Check(bytes.Equal(v1.([]byte), []byte(k.String())), Equals, true)
43 +
44 + v2, err := mpds.Get(ds.NewKey("abc").Child(k))
45 + c.Check(err, Equals, nil)
46 + c.Check(bytes.Equal(v2.([]byte), []byte(k.String())), Equals, true)
47 + }
48 +
49 + listA, errA := mpds.KeyList()
50 + listB, errB := nsds.KeyList()
51 + c.Check(errA, Equals, nil)
52 + c.Check(errB, Equals, nil)
53 + c.Check(len(listA), Equals, len(listB))
54 +
55 + // sort them cause yeah.
56 + sort.Sort(ds.KeySlice(listA))
57 + sort.Sort(ds.KeySlice(listB))
58 +
59 + for i, kA := range listA {
60 + kB := listB[i]
61 + c.Check(nsds.InvertKey(kA), Equals, kB)
62 + c.Check(kA, Equals, nsds.ConvertKey(kB))
63 + }
64 +}
65 +
66 +func strsToKeys(strs []string) []ds.Key {
67 + keys := make([]ds.Key, len(strs))
68 + for i, s := range strs {
69 + keys[i] = ds.NewKey(s)
70 + }
71 + return keys
72 +}
Godeps/_workspace/src/github.com/jbenet/go-datastore/fs/fs_test.go
+18 -4
@@ -2,6 +2,7 @@ package fs_test
2
3 import (
4 "bytes"
5 + "sort"
6 "testing"
7
8 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
@@ -12,10 +13,7 @@ import (
13 // Hook up gocheck into the "go test" runner.
14 func Test(t *testing.T) { TestingT(t) }
15
15 -type DSSuite struct {
16 - dir string
17 - ds ds.Datastore
18 -}
16 +type DSSuite struct{}
17
18 var _ = Suite(&DSSuite{})
19
@@ -54,6 +52,22 @@ func (ks *DSSuite) TestBasic(c *C) {
52 c.Check(err, Equals, nil)
53 c.Check(bytes.Equal(v.([]byte), []byte(k.String())), Equals, true)
54 }
55 +
56 + listA, errA := mpds.KeyList()
57 + listB, errB := ktds.KeyList()
58 + c.Check(errA, Equals, nil)
59 + c.Check(errB, Equals, nil)
60 + c.Check(len(listA), Equals, len(listB))
61 +
62 + // sort them cause yeah.
63 + sort.Sort(ds.KeySlice(listA))
64 + sort.Sort(ds.KeySlice(listB))
65 +
66 + for i, kA := range listA {
67 + kB := listB[i]
68 + c.Check(pair.Invert(kA), Equals, kB)
69 + c.Check(kA, Equals, pair.Convert(kB))
70 + }
71 }
72
73 func strsToKeys(strs []string) []ds.Key {
Godeps/_workspace/src/github.com/jbenet/go-datastore/io/io.go deleted
-44
@@ -1,44 +0,0 @@
1 -package leveldb
2 -
3 -import (
4 - "bytes"
5 - "io"
6 -
7 - ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
8 -)
9 -
10 -// CastAsReader does type assertions to find the type of a value and attempts
11 -// to turn it into an io.Reader. If not possible, will return ds.ErrInvalidType
12 -func CastAsReader(value interface{}) (io.Reader, error) {
13 - switch v := value.(type) {
14 - case io.Reader:
15 - return v, nil
16 -
17 - case []byte:
18 - return bytes.NewReader(v), nil
19 -
20 - case string:
21 - return bytes.NewReader([]byte(v)), nil
22 -
23 - default:
24 - return nil, ds.ErrInvalidType
25 - }
26 -}
27 -
28 -// // CastAsWriter does type assertions to find the type of a value and attempts
29 -// // to turn it into an io.Writer. If not possible, will return ds.ErrInvalidType
30 -// func CastAsWriter(value interface{}) (err error) {
31 -// switch v := value.(type) {
32 -// case io.Reader:
33 -// return v, nil
34 -//
35 -// case []byte:
36 -// return bytes.NewReader(v), nil
37 -//
38 -// case string:
39 -// return bytes.NewReader([]byte(v)), nil
40 -//
41 -// default:
42 -// return nil, ds.ErrInvalidType
43 -// }
44 -// }
blocks/set/dbset.go
+5 -7
@@ -9,19 +9,17 @@ import (
9 type datastoreBlockSet struct {
10 dstore ds.Datastore
11 bset BlockSet
12 - prefix string
12 }
13
15 -func NewDBWrapperSet(d ds.Datastore, prefix string, bset BlockSet) BlockSet {
14 +func NewDBWrapperSet(d ds.Datastore, bset BlockSet) BlockSet {
15 return &datastoreBlockSet{
16 dstore: d,
17 bset: bset,
19 - prefix: prefix,
18 }
19 }
20
21 func (d *datastoreBlockSet) AddBlock(k util.Key) {
24 - err := d.dstore.Put(d.prefixKey(k), []byte{})
22 + err := d.dstore.Put(k.DsKey(), []byte{})
23 if err != nil {
24 log.Error("blockset put error: %s", err)
25 }
@@ -32,7 +30,7 @@ func (d *datastoreBlockSet) AddBlock(k util.Key) {
30 func (d *datastoreBlockSet) RemoveBlock(k util.Key) {
31 d.bset.RemoveBlock(k)
32 if !d.bset.HasKey(k) {
35 - d.dstore.Delete(d.prefixKey(k))
33 + d.dstore.Delete(k.DsKey())
34 }
35 }
36
@@ -44,6 +42,6 @@ func (d *datastoreBlockSet) GetBloomFilter() bloom.Filter {
42 return d.bset.GetBloomFilter()
43 }
44
47 -func (d *datastoreBlockSet) prefixKey(k util.Key) ds.Key {
48 - return (util.Key(d.prefix) + k).DsKey()
45 +func (d *datastoreBlockSet) GetKeys() []util.Key {
46 + return d.bset.GetKeys()
47 }
blocks/set/set.go
+35
@@ -1,6 +1,10 @@
1 package set
2
3 import (
4 + "errors"
5 +
6 + ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
7 +
8 "github.com/jbenet/go-ipfs/blocks/bloom"
9 "github.com/jbenet/go-ipfs/util"
10 )
@@ -12,6 +16,29 @@ type BlockSet interface {
16 RemoveBlock(util.Key)
17 HasKey(util.Key) bool
18 GetBloomFilter() bloom.Filter
19 +
20 + GetKeys() []util.Key
21 +}
22 +
23 +func SimpleSetFromKeys(keys []util.Key) BlockSet {
24 + sbs := &simpleBlockSet{blocks: make(map[util.Key]struct{})}
25 + for _, k := range keys {
26 + sbs.blocks[k] = struct{}{}
27 + }
28 + return sbs
29 +}
30 +
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
44 func NewSimpleBlockSet() BlockSet {
@@ -42,3 +69,11 @@ func (b *simpleBlockSet) GetBloomFilter() bloom.Filter {
69 }
70 return f
71 }
72 +
73 +func (b *simpleBlockSet) GetKeys() []util.Key {
74 + var out []util.Key
75 + for k, _ := range b.blocks {
76 + out = append(out, k)
77 + }
78 + return out
79 +}
pin/indirect.go
+25 -9
@@ -1,25 +1,41 @@
1 package pin
2
3 import (
4 + "errors"
5 +
6 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
5 - bc "github.com/jbenet/go-ipfs/blocks/set"
7 + "github.com/jbenet/go-ipfs/blocks/set"
8 "github.com/jbenet/go-ipfs/util"
9 )
10
11 type indirectPin struct {
10 - blockset bc.BlockSet
12 + blockset set.BlockSet
13 refCounts map[util.Key]int
14 }
15
14 -func loadBlockSet(d ds.Datastore) (bc.BlockSet, map[util.Key]int) {
15 - panic("Not yet implemented!")
16 - return nil, nil
16 +func NewIndirectPin(dstore ds.Datastore) *indirectPin {
17 + return &indirectPin{
18 + blockset: set.NewDBWrapperSet(dstore, set.NewSimpleBlockSet()),
19 + refCounts: make(map[util.Key]int),
20 + }
21 }
22
19 -func newIndirectPin(d ds.Datastore) indirectPin {
20 - // suppose the blockset actually takes blocks, not just keys
21 - bs, rc := loadBlockSet(d)
22 - return indirectPin{bs, rc}
23 +func loadIndirPin(d ds.Datastore, k ds.Key) (*indirectPin, error) {
24 + irefcnt, err := d.Get(k)
25 + if err != nil {
26 + return nil, err
27 + }
28 + refcnt, ok := irefcnt.(map[util.Key]int)
29 + if !ok {
30 + return nil, errors.New("invalid type from datastore")
31 + }
32 +
33 + var keys []util.Key
34 + for k, _ := range refcnt {
35 + keys = append(keys, k)
36 + }
37 +
38 + return &indirectPin{blockset: set.SimpleSetFromKeys(keys), refCounts: refcnt}, nil
39 }
40
41 func (i *indirectPin) Increment(k util.Key) {
pin/pin.go
+57 -5
@@ -1,21 +1,29 @@
1 package pin
2
3 import (
4 +
5 + //ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
6 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
7 + nsds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go/namespace"
8 "github.com/jbenet/go-ipfs/blocks/set"
9 mdag "github.com/jbenet/go-ipfs/merkledag"
10 "github.com/jbenet/go-ipfs/util"
11 )
12
13 +var recursePinDatastoreKey = ds.NewKey("/local/pins/recursive/keys")
14 +var directPinDatastoreKey = ds.NewKey("/local/pins/direct/keys")
15 +var indirectPinDatastoreKey = ds.NewKey("/local/pins/indirect/keys")
16 +
17 type Pinner interface {
18 Pin(*mdag.Node, bool) error
19 Unpin(util.Key, bool) error
20 + Flush() error
21 }
22
23 type pinner struct {
24 recursePin set.BlockSet
25 directPin set.BlockSet
18 - indirPin indirectPin
26 + indirPin *indirectPin
27 dserv *mdag.DAGService
28 dstore ds.Datastore
29 }
@@ -23,14 +31,17 @@ type pinner struct {
31 func NewPinner(dstore ds.Datastore, serv *mdag.DAGService) Pinner {
32
33 // Load set from given datastore...
26 - rcset := set.NewDBWrapperSet(dstore, "/pinned/recurse/", set.NewSimpleBlockSet())
27 - dirset := set.NewDBWrapperSet(dstore, "/pinned/direct/", set.NewSimpleBlockSet())
34 + rcds := nsds.Wrap(dstore, recursePinDatastoreKey)
35 + rcset := set.NewDBWrapperSet(rcds, set.NewSimpleBlockSet())
36
29 - nsdstore := dstore // WRAP IN NAMESPACE
37 + dirds := nsds.Wrap(dstore, directPinDatastoreKey)
38 + dirset := set.NewDBWrapperSet(dirds, set.NewSimpleBlockSet())
39 +
40 + nsdstore := nsds.Wrap(dstore, indirectPinDatastoreKey)
41 return &pinner{
42 recursePin: rcset,
43 directPin: dirset,
33 - indirPin: newIndirectPin(nsdstore),
44 + indirPin: NewIndirectPin(nsdstore),
45 dserv: serv,
46 dstore: dstore,
47 }
@@ -126,3 +137,44 @@ func (p *pinner) IsPinned(key util.Key) bool {
137 p.directPin.HasKey(key) ||
138 p.indirPin.HasKey(key)
139 }
140 +
141 +func LoadPinner(d ds.Datastore) (Pinner, error) {
142 + p := new(pinner)
143 +
144 + var err error
145 + p.recursePin, err = set.SetFromDatastore(d, recursePinDatastoreKey)
146 + if err != nil {
147 + return nil, err
148 + }
149 + p.directPin, err = set.SetFromDatastore(d, directPinDatastoreKey)
150 + if err != nil {
151 + return nil, err
152 + }
153 +
154 + p.indirPin, err = loadIndirPin(d, indirectPinDatastoreKey)
155 + if err != nil {
156 + return nil, err
157 + }
158 +
159 + return p, nil
160 +}
161 +
162 +func (p *pinner) Flush() error {
163 + recurse := p.recursePin.GetKeys()
164 + err := p.dstore.Put(recursePinDatastoreKey, recurse)
165 + if err != nil {
166 + return err
167 + }
168 +
169 + direct := p.directPin.GetKeys()
170 + err = p.dstore.Put(directPinDatastoreKey, direct)
171 + if err != nil {
172 + return err
173 + }
174 +
175 + err = p.dstore.Put(indirectPinDatastoreKey, p.indirPin.refCounts)
176 + if err != nil {
177 + return err
178 + }
179 + return nil
180 +}