@cryptotaxi247 / kubo / commits / 275b03f81

rework dagreader to have a dagservice for node resolution

Jeromy committed Sep 6, 2014 at 22:11 UTC 275b03f814b784adf72468425a3b2a9cdd54f283
6 files changed +66 -29
cmd/ipfs/cat.go
+1 -23
@@ -7,7 +7,6 @@ import (
7
8 "github.com/gonuts/flag"
9 "github.com/jbenet/commander"
10 - bserv "github.com/jbenet/go-ipfs/blockservice"
10 dag "github.com/jbenet/go-ipfs/merkledag"
11 u "github.com/jbenet/go-ipfs/util"
12 )
@@ -41,12 +40,7 @@ func catCmd(c *commander.Command, inp []string) error {
40 return err
41 }
42
44 - err = ExpandDag(nd, n.Blocks)
45 - if err != nil {
46 - return err
47 - }
48 -
49 - read, err := dag.NewDagReader(nd)
43 + read, err := dag.NewDagReader(nd, n.DAG)
44 if err != nil {
45 fmt.Println(err)
46 continue
@@ -60,19 +54,3 @@ func catCmd(c *commander.Command, inp []string) error {
54 }
55 return nil
56 }
63 -
64 -// Expand all subnodes in this dag so printing can occur without error
65 -//TODO: this needs to be done MUCH better in a somewhat asynchronous way.
66 -//also should be moved elsewhere.
67 -func ExpandDag(nd *dag.Node, bs *bserv.BlockService) error {
68 - for _, lnk := range nd.Links {
69 - if lnk.Node == nil {
70 - blk, err := bs.GetBlock(u.Key(lnk.Hash))
71 - if err != nil {
72 - return err
73 - }
74 - lnk.Node = &dag.Node{Data: dag.WrapData(blk.Data)}
75 - }
76 - }
77 - return nil
78 -}
fuse/readonly/readonly_unix.go
+2 -1
@@ -74,6 +74,7 @@ func (*Root) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
74 type Node struct {
75 Ipfs *core.IpfsNode
76 Nd *mdag.Node
77 + fd *mdag.DagReader
78 }
79
80 // Attr returns the attributes of a given node.
@@ -120,7 +121,7 @@ func (s *Node) ReadDir(intr fs.Intr) ([]fuse.Dirent, fuse.Error) {
121 // ReadAll reads the object data as file data
122 func (s *Node) ReadAll(intr fs.Intr) ([]byte, fuse.Error) {
123 u.DOut("Read node.\n")
123 - r, err := mdag.NewDagReader(s.Nd)
124 + r, err := mdag.NewDagReader(s.Nd, s.Ipfs.DAG)
125 if err != nil {
126 return nil, err
127 }
importer/importer_test.go
+24
@@ -32,3 +32,27 @@ func TestFileConsistency(t *testing.T) {
32 t.Fatal("Output not the same as input.")
33 }
34 }
35 +
36 +//Test where calls to read are smaller than the chunk size
37 +func TestFileConsistencyLargeBlocks(t *testing.T) {
38 + buf := new(bytes.Buffer)
39 + io.CopyN(buf, rand.Reader, 4096*32)
40 + should := buf.Bytes()
41 + nd, err := NewDagFromReaderWithSplitter(buf, SplitterBySize(4096))
42 + if err != nil {
43 + t.Fatal(err)
44 + }
45 + r, err := dag.NewDagReader(nd)
46 + if err != nil {
47 + t.Fatal(err)
48 + }
49 +
50 + out, err := ioutil.ReadAll(r)
51 + if err != nil {
52 + t.Fatal(err)
53 + }
54 +
55 + if !bytes.Equal(out, should) {
56 + t.Fatal("Output not the same as input.")
57 + }
58 +}
merkledag/dagreader.go
+36 -4
@@ -5,20 +5,22 @@ import (
5 "errors"
6 "io"
7
8 - "code.google.com/p/goprotobuf/proto"
8 + proto "code.google.com/p/goprotobuf/proto"
9 + u "github.com/jbenet/go-ipfs/util"
10 )
11
12 var ErrIsDir = errors.New("this dag node is a directory.")
13
14 // DagReader provides a way to easily read the data contained in a dag.
15 type DagReader struct {
16 + serv *DAGService
17 node *Node
18 position int
19 buf *bytes.Buffer
20 thisData []byte
21 }
22
21 -func NewDagReader(n *Node) (io.Reader, error) {
23 +func NewDagReader(n *Node, serv *DAGService) (io.Reader, error) {
24 pb := new(PBData)
25 err := proto.Unmarshal(n.Data, pb)
26 if err != nil {
@@ -31,6 +33,7 @@ func NewDagReader(n *Node) (io.Reader, error) {
33 return &DagReader{
34 node: n,
35 thisData: pb.GetData(),
36 + serv: serv,
37 }, nil
38 case PBData_Raw:
39 return bytes.NewBuffer(pb.GetData()), nil
@@ -46,8 +49,11 @@ func (dr *DagReader) precalcNextBuf() error {
49 nxtLink := dr.node.Links[dr.position]
50 nxt := nxtLink.Node
51 if nxt == nil {
49 - //TODO: should use dagservice or something to get needed block
50 - return errors.New("Link to nil node! Tree not fully expanded!")
52 + nxtNode, err := dr.serv.Get(u.Key(nxtLink.Hash))
53 + if err != nil {
54 + return err
55 + }
56 + nxt = nxtNode
57 }
58 pb := new(PBData)
59 err := proto.Unmarshal(nxt.Data, pb)
@@ -96,3 +102,29 @@ func (dr *DagReader) Read(b []byte) (int, error) {
102 }
103 }
104 }
105 +
106 +/*
107 +func (dr *DagReader) Seek(offset int64, whence int) (int64, error) {
108 + switch whence {
109 + case os.SEEK_SET:
110 + for i := 0; i < len(dr.node.Links); i++ {
111 + nsize := dr.node.Links[i].Size - 8
112 + if offset > nsize {
113 + offset -= nsize
114 + } else {
115 + break
116 + }
117 + }
118 + dr.position = i
119 + err := dr.precalcNextBuf()
120 + if err != nil {
121 + return 0, err
122 + }
123 + case os.SEEK_CUR:
124 + case os.SEEK_END:
125 + default:
126 + return 0, errors.New("invalid whence")
127 + }
128 + return 0, nil
129 +}
130 +*/
merkledag/merkledag.go
+2
@@ -96,6 +96,8 @@ func (n *Node) Key() (u.Key, error) {
96 // DAGService is an IPFS Merkle DAG service.
97 // - the root is virtual (like a forest)
98 // - stores nodes' data in a BlockService
99 +// TODO: should cache Nodes that are in memory, and be
100 +// able to free some of them when vm pressure is high
101 type DAGService struct {
102 Blocks *bserv.BlockService
103 }
routing/dht/routing.go
+1 -1
@@ -192,7 +192,7 @@ func (dht *IpfsDHT) FindProvidersAsync(key u.Key, count int, timeout time.Durati
192 provs := dht.providers.GetProviders(key)
193 for _, p := range provs {
194 count--
195 - // NOTE: assuming that the list of peers is unique
195 + // NOTE: assuming that this list of peers is unique
196 ps.Add(p)
197 peerOut <- p
198 if count <= 0 {