@cryptotaxi247 / kubo / commits / 056699ceb

convert DAGService to an interface

Emery Hemingway committed Oct 25, 2014 at 22:15 UTC 056699cebe09afb362c731fc94964d049d81e557
11 files changed +91 -32
core/core.go
+2 -2
@@ -58,7 +58,7 @@ type IpfsNode struct {
58 Blocks *bserv.BlockService
59
60 // the merkle dag service, get/add objects.
61 - DAG *merkledag.DAGService
61 + DAG merkledag.DAGService
62
63 // the path resolution system
64 Resolver *path.Resolver
@@ -161,7 +161,7 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
161 return nil, err
162 }
163
164 - dag := &merkledag.DAGService{Blocks: bs}
164 + dag := merkledag.NewDAGService(bs)
165 ns := namesys.NewNameSystem(route)
166 p, err := pin.LoadPinner(d, dag)
167 if err != nil {
core/mock.go
+1 -1
@@ -49,7 +49,7 @@ func NewMockNode() (*IpfsNode, error) {
49 return nil, err
50 }
51
52 - nd.DAG = &mdag.DAGService{Blocks: bserv}
52 + nd.DAG = mdag.NewDAGService(bserv)
53
54 // Namespace resolver
55 nd.Namesys = nsys.NewNameSystem(dht)
merkledag/merkledag.go
+22 -10
@@ -62,7 +62,7 @@ func MakeLink(n *Node) (*Link, error) {
62 }, nil
63 }
64
65 -func (l *Link) GetNode(serv *DAGService) (*Node, error) {
65 +func (l *Link) GetNode(serv DAGService) (*Node, error) {
66 if l.Node != nil {
67 return l.Node, nil
68 }
@@ -151,20 +151,32 @@ func (n *Node) Key() (u.Key, error) {
151 }
152
153 // DAGService is an IPFS Merkle DAG service.
154 +type DAGService interface {
155 + Add(*Node) (u.Key, error)
156 + AddRecursive(*Node) error
157 + Get(u.Key) (*Node, error)
158 + Remove(*Node) error
159 +}
160 +
161 +func NewDAGService(bs *bserv.BlockService) DAGService {
162 + return &dagService{bs}
163 +}
164 +
165 +// dagService is an IPFS Merkle DAG service.
166 // - the root is virtual (like a forest)
167 // - stores nodes' data in a BlockService
168 // TODO: should cache Nodes that are in memory, and be
169 // able to free some of them when vm pressure is high
158 -type DAGService struct {
170 +type dagService struct {
171 Blocks *bserv.BlockService
172 }
173
162 -// Add adds a node to the DAGService, storing the block in the BlockService
163 -func (n *DAGService) Add(nd *Node) (u.Key, error) {
174 +// Add adds a node to the dagService, storing the block in the BlockService
175 +func (n *dagService) Add(nd *Node) (u.Key, error) {
176 k, _ := nd.Key()
177 log.Debug("DagService Add [%s]", k)
178 if n == nil {
167 - return "", fmt.Errorf("DAGService is nil")
179 + return "", fmt.Errorf("dagService is nil")
180 }
181
182 d, err := nd.Encoded(false)
@@ -182,7 +194,7 @@ func (n *DAGService) Add(nd *Node) (u.Key, error) {
194 return n.Blocks.AddBlock(b)
195 }
196
185 -func (n *DAGService) AddRecursive(nd *Node) error {
197 +func (n *dagService) AddRecursive(nd *Node) error {
198 _, err := n.Add(nd)
199 if err != nil {
200 log.Info("AddRecursive Error: %s\n", err)
@@ -201,10 +213,10 @@ func (n *DAGService) AddRecursive(nd *Node) error {
213 return nil
214 }
215
204 -// Get retrieves a node from the DAGService, fetching the block in the BlockService
205 -func (n *DAGService) Get(k u.Key) (*Node, error) {
216 +// Get retrieves a node from the dagService, fetching the block in the BlockService
217 +func (n *dagService) Get(k u.Key) (*Node, error) {
218 if n == nil {
207 - return nil, fmt.Errorf("DAGService is nil")
219 + return nil, fmt.Errorf("dagService is nil")
220 }
221
222 ctx, _ := context.WithTimeout(context.TODO(), time.Second*5)
@@ -216,7 +228,7 @@ func (n *DAGService) Get(k u.Key) (*Node, error) {
228 return Decoded(b.Data)
229 }
230
219 -func (n *DAGService) Remove(nd *Node) error {
231 +func (n *dagService) Remove(nd *Node) error {
232 for _, l := range nd.Links {
233 if l.Node != nil {
234 n.Remove(l.Node)
path/path.go
+1 -1
@@ -15,7 +15,7 @@ var log = u.Logger("path")
15 // Resolver provides path resolution to IPFS
16 // It has a pointer to a DAGService, which is uses to resolve nodes.
17 type Resolver struct {
18 - DAG *merkledag.DAGService
18 + DAG merkledag.DAGService
19 }
20
21 // ResolvePath fetches the node for given path. It uses the first
pin/pin.go
+3 -3
@@ -32,11 +32,11 @@ type pinner struct {
32 recursePin set.BlockSet
33 directPin set.BlockSet
34 indirPin *indirectPin
35 - dserv *mdag.DAGService
35 + dserv mdag.DAGService
36 dstore ds.Datastore
37 }
38
39 -func NewPinner(dstore ds.Datastore, serv *mdag.DAGService) Pinner {
39 +func NewPinner(dstore ds.Datastore, serv mdag.DAGService) Pinner {
40
41 // Load set from given datastore...
42 rcds := nsds.Wrap(dstore, recursePinDatastoreKey)
@@ -151,7 +151,7 @@ func (p *pinner) IsPinned(key util.Key) bool {
151 p.indirPin.HasKey(key)
152 }
153
154 -func LoadPinner(d ds.Datastore, dserv *mdag.DAGService) (Pinner, error) {
154 +func LoadPinner(d ds.Datastore, dserv mdag.DAGService) (Pinner, error) {
155 p := new(pinner)
156
157 { // load recursive set
pin/pin_test.go
+1 -1
@@ -24,7 +24,7 @@ func TestPinnerBasic(t *testing.T) {
24 t.Fatal(err)
25 }
26
27 - dserv := &mdag.DAGService{Blocks: bserv}
27 + dserv := mdag.NewDAGService(bserv)
28
29 p := NewPinner(dstore, dserv)
30
unixfs/io/dagmodifier.go
+2 -2
@@ -17,14 +17,14 @@ import (
17 // perform surgery on a DAG 'file'
18 // Dear god, please rename this to something more pleasant
19 type DagModifier struct {
20 - dagserv *mdag.DAGService
20 + dagserv mdag.DAGService
21 curNode *mdag.Node
22
23 pbdata *ftpb.Data
24 splitter chunk.BlockSplitter
25 }
26
27 -func NewDagModifier(from *mdag.Node, serv *mdag.DAGService, spl chunk.BlockSplitter) (*DagModifier, error) {
27 +func NewDagModifier(from *mdag.Node, serv mdag.DAGService, spl chunk.BlockSplitter) (*DagModifier, error) {
28 pbd, err := ft.FromBytes(from.Data)
29 if err != nil {
30 return nil, err
unixfs/io/dagmodifier_test.go
+9 -5
@@ -16,17 +16,17 @@ import (
16 logging "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-logging"
17 )
18
19 -func getMockDagServ(t *testing.T) *mdag.DAGService {
19 +func getMockDagServ(t *testing.T) mdag.DAGService {
20 dstore := ds.NewMapDatastore()
21 bserv, err := bs.NewBlockService(dstore, nil)
22 if err != nil {
23 t.Fatal(err)
24 }
25 - return &mdag.DAGService{Blocks: bserv}
25 + return mdag.NewDAGService(bserv)
26 }
27
28 -func getNode(t *testing.T, dserv *mdag.DAGService, size int64) ([]byte, *mdag.Node) {
29 - dw := NewDagWriter(dserv, &chunk.SizeSplitter{Size: 500})
28 +func getNode(t *testing.T, dserv mdag.DAGService, size int64) ([]byte, *mdag.Node) {
29 + dw := NewDagWriter(dserv, &chunk.SizeSplitter{500})
30
31 n, err := io.CopyN(dw, u.NewFastRand(), size)
32 if err != nil {
@@ -36,7 +36,11 @@ func getNode(t *testing.T, dserv *mdag.DAGService, size int64) ([]byte, *mdag.No
36 t.Fatal("Incorrect copy amount!")
37 }
38
39 - dw.Close()
39 + err = dw.Close()
40 + if err != nil {
41 + t.Fatal("DagWriter failed to close,", err)
42 + }
43 +
44 node := dw.GetNode()
45
46 dr, err := NewDagReader(node, dserv)
unixfs/io/dagreader.go
+2 -2
@@ -16,7 +16,7 @@ var ErrIsDir = errors.New("this dag node is a directory")
16
17 // DagReader provides a way to easily read the data contained in a dag.
18 type DagReader struct {
19 - serv *mdag.DAGService
19 + serv mdag.DAGService
20 node *mdag.Node
21 position int
22 buf *bytes.Buffer
@@ -24,7 +24,7 @@ type DagReader struct {
24
25 // NewDagReader creates a new reader object that reads the data represented by the given
26 // node, using the passed in DAGService for data retreival
27 -func NewDagReader(n *mdag.Node, serv *mdag.DAGService) (io.Reader, error) {
27 +func NewDagReader(n *mdag.Node, serv mdag.DAGService) (io.Reader, error) {
28 pb := new(ftpb.Data)
29 err := proto.Unmarshal(n.Data, pb)
30 if err != nil {
unixfs/io/dagwriter.go
+2 -2
@@ -10,7 +10,7 @@ import (
10 var log = util.Logger("dagwriter")
11
12 type DagWriter struct {
13 - dagserv *dag.DAGService
13 + dagserv dag.DAGService
14 node *dag.Node
15 totalSize int64
16 splChan chan []byte
@@ -19,7 +19,7 @@ type DagWriter struct {
19 seterr error
20 }
21
22 -func NewDagWriter(ds *dag.DAGService, splitter chunk.BlockSplitter) *DagWriter {
22 +func NewDagWriter(ds dag.DAGService, splitter chunk.BlockSplitter) *DagWriter {
23 dw := new(DagWriter)
24 dw.dagserv = ds
25 dw.splChan = make(chan []byte, 8)
unixfs/io/dagwriter_test.go
+46 -3
@@ -7,6 +7,7 @@ import (
7
8 ds "github.com/jbenet/go-ipfs/Godeps/_workspace/src/github.com/jbenet/go-datastore"
9 bs "github.com/jbenet/go-ipfs/blockservice"
10 + "github.com/jbenet/go-ipfs/importer"
11 chunk "github.com/jbenet/go-ipfs/importer/chunk"
12 mdag "github.com/jbenet/go-ipfs/merkledag"
13 )
@@ -53,7 +54,7 @@ func TestDagWriter(t *testing.T) {
54 if err != nil {
55 t.Fatal(err)
56 }
56 - dag := &mdag.DAGService{Blocks: bserv}
57 + dag := mdag.NewDAGService(bserv)
58 dw := NewDagWriter(dag, &chunk.SizeSplitter{Size: 4096})
59
60 nbytes := int64(1024 * 1024 * 2)
@@ -87,7 +88,7 @@ func TestMassiveWrite(t *testing.T) {
88 if err != nil {
89 t.Fatal(err)
90 }
90 - dag := &mdag.DAGService{Blocks: bserv}
91 + dag := mdag.NewDAGService(bserv)
92 dw := NewDagWriter(dag, &chunk.SizeSplitter{Size: 4096})
93
94 nbytes := int64(1024 * 1024 * 1024 * 16)
@@ -107,7 +108,7 @@ func BenchmarkDagWriter(b *testing.B) {
108 if err != nil {
109 b.Fatal(err)
110 }
110 - dag := &mdag.DAGService{Blocks: bserv}
111 + dag := mdag.NewDAGService(bserv)
112
113 b.ResetTimer()
114 nbytes := int64(100000)
@@ -125,3 +126,45 @@ func BenchmarkDagWriter(b *testing.B) {
126 }
127
128 }
129 +
130 +func TestAgainstImporter(t *testing.T) {
131 + dstore := ds.NewMapDatastore()
132 + bserv, err := bs.NewBlockService(dstore, nil)
133 + if err != nil {
134 + t.Fatal(err)
135 + }
136 + dag := mdag.NewDAGService(bserv)
137 +
138 + nbytes := int64(1024 * 1024 * 2)
139 +
140 + // DagWriter
141 + dw := NewDagWriter(dag, &chunk.SizeSplitter{4096})
142 + n, err := io.CopyN(dw, &datasource{}, nbytes)
143 + if err != nil {
144 + t.Fatal(err)
145 + }
146 + if n != nbytes {
147 + t.Fatal("Copied incorrect amount of bytes!")
148 + }
149 +
150 + dw.Close()
151 + dwNode := dw.GetNode()
152 + dwKey, err := dwNode.Key()
153 + if err != nil {
154 + t.Fatal(err)
155 + }
156 +
157 + // DagFromFile
158 + rl := &io.LimitedReader{&datasource{}, nbytes}
159 +
160 + dffNode, err := importer.NewDagFromReaderWithSplitter(rl, &chunk.SizeSplitter{4096})
161 + dffKey, err := dffNode.Key()
162 + if err != nil {
163 + t.Fatal(err)
164 + }
165 + if dwKey.String() != dffKey.String() {
166 + t.Errorf("\nDagWriter produced %s\n"+
167 + "DagFromReader produced %s",
168 + dwKey, dffKey)
169 + }
170 +}