implement unpin and add a datastore backed blockset
Jeromy committed
Oct 16, 2014 at 19:20 UTC
51879e96402937fe57caed875dc55c50671241a4
5 files changed
+119
-7
blocks/bloom/filter.go
+25
@@ -1,6 +1,7 @@
1
package bloom
2
3
import (
4
+ "errors"
5
"fmt"
6
"hash"
7
"hash/adler32"
@@ -12,6 +13,7 @@ import (
13
type Filter interface {
14
Add([]byte)
15
Find([]byte) bool
16
+ Merge(Filter) (Filter, error)
17
}
18
19
func BasicFilter() Filter {
@@ -68,3 +70,26 @@ func bytesMod(b []byte, modulo int64) int64 {
70
71
return result.Int64()
72
}
73
+
74
+func (f *filter) Merge(o Filter) (Filter, error) {
75
+ casfil, ok := o.(*filter)
76
+ if !ok {
77
+ return nil, errors.New("Unsupported filter type")
78
+ }
79
+
80
+ if len(casfil.filter) != len(f.filter) {
81
+ return nil, errors.New("filter lengths must match!")
82
+ }
83
+
84
+ nfilt := new(filter)
85
+
86
+ // this bit is sketchy, need a way of comparing hash functions
87
+ nfilt.hashes = f.hashes
88
+
89
+ nfilt.filter = make([]byte, len(f.filter))
90
+ for i, v := range f.filter {
91
+ nfilt.filter[i] = v | casfil.filter[i]
92
+ }
93
+
94
+ return nfilt, nil
95
+}
blocks/set/dbset.go
new
+49
@@ -0,0 +1,49 @@
1
+package set
2
+
3
+import (
4
+ ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/datastore.go"
5
+ "github.com/jbenet/go-ipfs/blocks/bloom"
6
+ "github.com/jbenet/go-ipfs/util"
7
+)
8
+
9
+type datastoreBlockSet struct {
10
+ dstore ds.Datastore
11
+ bset BlockSet
12
+ prefix string
13
+}
14
+
15
+func NewDBWrapperSet(d ds.Datastore, prefix string, bset BlockSet) BlockSet {
16
+ return &datastoreBlockSet{
17
+ dstore: d,
18
+ bset: bset,
19
+ prefix: prefix,
20
+ }
21
+}
22
+
23
+func (d *datastoreBlockSet) AddBlock(k util.Key) {
24
+ err := d.dstore.Put(d.prefixKey(k), []byte{})
25
+ if err != nil {
26
+ log.Error("blockset put error: %s", err)
27
+ }
28
+
29
+ d.bset.AddBlock(k)
30
+}
31
+
32
+func (d *datastoreBlockSet) RemoveBlock(k util.Key) {
33
+ d.bset.RemoveBlock(k)
34
+ if !d.bset.HasKey(k) {
35
+ d.dstore.Delete(d.prefixKey(k))
36
+ }
37
+}
38
+
39
+func (d *datastoreBlockSet) HasKey(k util.Key) bool {
40
+ return d.bset.HasKey(k)
41
+}
42
+
43
+func (d *datastoreBlockSet) GetBloomFilter() bloom.Filter {
44
+ return d.bset.GetBloomFilter()
45
+}
46
+
47
+func (d *datastoreBlockSet) prefixKey(k util.Key) ds.Key {
48
+ return (util.Key(d.prefix) + k).DsKey()
49
+}
blocks/set/refset.go
renamed
+2
-3
@@ -1,8 +1,7 @@
1
-package pin
1
+package set
2
3
import (
4
"github.com/jbenet/go-ipfs/blocks/bloom"
5
- "github.com/jbenet/go-ipfs/blocks/set"
5
"github.com/jbenet/go-ipfs/util"
6
)
7
@@ -10,7 +9,7 @@ type refCntBlockSet struct {
9
blocks map[util.Key]int
10
}
11
13
-func NewRefCountBlockSet() set.BlockSet {
12
+func NewRefCountBlockSet() BlockSet {
13
return &refCntBlockSet{blocks: make(map[util.Key]int)}
14
}
15
blocks/set/set.go
+2
@@ -5,6 +5,8 @@ import (
5
"github.com/jbenet/go-ipfs/util"
6
)
7
8
+var log = util.Logger("blockset")
9
+
10
type BlockSet interface {
11
AddBlock(util.Key)
12
RemoveBlock(util.Key)
pin/pin.go
+41
-4
@@ -17,14 +17,19 @@ type pinner struct {
17
directPin set.BlockSet
18
indirPin set.BlockSet
19
dserv *mdag.DAGService
20
+ dstore ds.Datastore
21
}
22
23
func NewPinner(dstore ds.Datastore, serv *mdag.DAGService) Pinner {
24
+ rcset := set.NewDBWrapperSet(dstore, "/pinned/recurse/", set.NewSimpleBlockSet())
25
+ dirset := set.NewDBWrapperSet(dstore, "/pinned/direct/", set.NewSimpleBlockSet())
26
+ indset := set.NewDBWrapperSet(dstore, "/pinned/indirect/", set.NewRefCountBlockSet())
27
return &pinner{
24
- recursePin: set.NewSimpleBlockSet(),
25
- directPin: set.NewSimpleBlockSet(),
26
- indirPin: NewRefCountBlockSet(),
28
+ recursePin: rcset,
29
+ directPin: dirset,
30
+ indirPin: indset,
31
dserv: serv,
32
+ dstore: dstore,
33
}
34
}
35
@@ -52,7 +57,39 @@ func (p *pinner) Pin(node *mdag.Node, recurse bool) error {
57
}
58
59
func (p *pinner) Unpin(k util.Key, recurse bool) error {
55
- panic("not yet implemented!")
60
+ if recurse {
61
+ p.recursePin.RemoveBlock(k)
62
+ node, err := p.dserv.Get(k)
63
+ if err != nil {
64
+ return err
65
+ }
66
+
67
+ return p.unpinLinks(node)
68
+ } else {
69
+ p.directPin.RemoveBlock(k)
70
+ }
71
+ return nil
72
+}
73
+
74
+func (p *pinner) unpinLinks(node *mdag.Node) error {
75
+ for _, l := range node.Links {
76
+ node, err := l.GetNode(p.dserv)
77
+ if err != nil {
78
+ return err
79
+ }
80
+
81
+ k, err := node.Key()
82
+ if err != nil {
83
+ return err
84
+ }
85
+
86
+ p.recursePin.RemoveBlock(k)
87
+
88
+ err = p.unpinLinks(node)
89
+ if err != nil {
90
+ return err
91
+ }
92
+ }
93
return nil
94
}
95