@cryptotaxi247 / kubo / commits / dfa0351df

Refactor ipfs get

License: MIT Signed-off-by: rht <rhtbot@gmail.com>

rht committed Aug 10, 2015 at 04:05 UTC dfa0351df93042973774ea66dd19aee487d7d77d
4 files changed +70 -103
core/commands/get.go
+8 -49
@@ -1,7 +1,6 @@
1 package commands
2
3 import (
4 - "bufio"
4 "compress/gzip"
5 "errors"
6 "fmt"
@@ -11,13 +10,11 @@ import (
10 "strings"
11
12 "github.com/ipfs/go-ipfs/Godeps/_workspace/src/github.com/cheggaaa/pb"
14 - context "github.com/ipfs/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
13
14 cmds "github.com/ipfs/go-ipfs/commands"
15 core "github.com/ipfs/go-ipfs/core"
16 path "github.com/ipfs/go-ipfs/path"
17 tar "github.com/ipfs/go-ipfs/thirdparty/tar"
20 - uio "github.com/ipfs/go-ipfs/unixfs/io"
18 utar "github.com/ipfs/go-ipfs/unixfs/tar"
19 )
20
@@ -64,15 +61,16 @@ may also specify the level of compression by specifying '-l=<1-9>'.
61 res.SetError(err, cmds.ErrNormal)
62 return
63 }
67 -
64 p := path.Path(req.Arguments()[0])
69 - var reader io.Reader
70 - if archive, _, _ := req.Option("archive").Bool(); !archive && cmplvl != gzip.NoCompression {
71 - // only use this when the flag is '-C' without '-a'
72 - reader, err = getZip(req.Context(), node, p, cmplvl)
73 - } else {
74 - reader, err = get(req.Context(), node, p, cmplvl)
65 + ctx := req.Context()
66 + dn, err := core.Resolve(ctx, node, p)
67 + if err != nil {
68 + res.SetError(err, cmds.ErrNormal)
69 + return
70 }
71 +
72 + archive, _, _ := req.Option("archive").Bool()
73 + reader, err := utar.DagArchive(ctx, dn, p.String(), node.DAG, archive, cmplvl)
74 if err != nil {
75 res.SetError(err, cmds.ErrNormal)
76 return
@@ -192,42 +190,3 @@ func getCompressOptions(req cmds.Request) (int, error) {
190 }
191 return gzip.NoCompression, nil
192 }
195 -
196 -func get(ctx context.Context, node *core.IpfsNode, p path.Path, compression int) (io.Reader, error) {
197 - dn, err := core.Resolve(ctx, node, p)
198 - if err != nil {
199 - return nil, err
200 - }
201 -
202 - return utar.DagArchive(ctx, dn, p.String(), node.DAG, compression)
203 -}
204 -
205 -// getZip is equivalent to `ipfs getdag $hash | gzip`
206 -func getZip(ctx context.Context, node *core.IpfsNode, p path.Path, compression int) (io.Reader, error) {
207 - dagnode, err := core.Resolve(ctx, node, p)
208 - if err != nil {
209 - return nil, err
210 - }
211 -
212 - reader, err := uio.NewDagReader(ctx, dagnode, node.DAG)
213 - if err != nil {
214 - return nil, err
215 - }
216 -
217 - pr, pw := io.Pipe()
218 - gw, err := gzip.NewWriterLevel(pw, compression)
219 - if err != nil {
220 - return nil, err
221 - }
222 - bufin := bufio.NewReader(reader)
223 - go func() {
224 - _, err := bufin.WriteTo(gw)
225 - if err != nil {
226 - log.Error("Fail to compress the stream")
227 - }
228 - gw.Close()
229 - pw.Close()
230 - }()
231 -
232 - return pr, nil
233 -}
fuse/readonly/readonly_unix.go
+1 -1
@@ -180,7 +180,7 @@ func (s *Node) Read(ctx context.Context, req *fuse.ReadRequest, resp *fuse.ReadR
180 return err
181 }
182
183 - buf := resp.Data[:min(req.Size, int(r.Size()-req.Offset))]
183 + buf := resp.Data[:min(req.Size, int(int64(r.Size())-req.Offset))]
184 n, err := io.ReadFull(r, buf)
185 if err != nil && err != io.EOF {
186 return err
unixfs/io/dagreader.go
+6 -7
@@ -58,8 +58,7 @@ type ReadSeekCloser interface {
58 // node, using the passed in DAGService for data retreival
59 func NewDagReader(ctx context.Context, n *mdag.Node, serv mdag.DAGService) (*DagReader, error) {
60 pb := new(ftpb.Data)
61 - err := proto.Unmarshal(n.Data, pb)
62 - if err != nil {
61 + if err := proto.Unmarshal(n.Data, pb); err != nil {
62 return nil, err
63 }
64
@@ -70,7 +69,7 @@ func NewDagReader(ctx context.Context, n *mdag.Node, serv mdag.DAGService) (*Dag
69 case ftpb.Data_Raw:
70 fallthrough
71 case ftpb.Data_File:
73 - return newDataFileReader(ctx, n, pb, serv), nil
72 + return NewDataFileReader(ctx, n, pb, serv), nil
73 case ftpb.Data_Metadata:
74 if len(n.Links) == 0 {
75 return nil, errors.New("incorrectly formatted metadata object")
@@ -85,7 +84,7 @@ func NewDagReader(ctx context.Context, n *mdag.Node, serv mdag.DAGService) (*Dag
84 }
85 }
86
88 -func newDataFileReader(ctx context.Context, n *mdag.Node, pb *ftpb.Data, serv mdag.DAGService) *DagReader {
87 +func NewDataFileReader(ctx context.Context, n *mdag.Node, pb *ftpb.Data, serv mdag.DAGService) *DagReader {
88 fctx, cancel := context.WithCancel(ctx)
89 promises := serv.GetDAG(fctx, n)
90 return &DagReader{
@@ -124,7 +123,7 @@ func (dr *DagReader) precalcNextBuf(ctx context.Context) error {
123 // A directory should not exist within a file
124 return ft.ErrInvalidDirLocation
125 case ftpb.Data_File:
127 - dr.buf = newDataFileReader(dr.ctx, nxt, pb, dr.serv)
126 + dr.buf = NewDataFileReader(dr.ctx, nxt, pb, dr.serv)
127 return nil
128 case ftpb.Data_Raw:
129 dr.buf = NewRSNCFromBytes(pb.GetData())
@@ -137,8 +136,8 @@ func (dr *DagReader) precalcNextBuf(ctx context.Context) error {
136 }
137
138 // Size return the total length of the data from the DAG structured file.
140 -func (dr *DagReader) Size() int64 {
141 - return int64(dr.pbdata.GetFilesize())
139 +func (dr *DagReader) Size() uint64 {
140 + return dr.pbdata.GetFilesize()
141 }
142
143 // Read reads data from the DAG structured file
unixfs/tar/writer.go
+55 -46
@@ -4,7 +4,6 @@ import (
4 "archive/tar"
5 "bufio"
6 "compress/gzip"
7 - "fmt"
7 "io"
8 "path"
9 "time"
@@ -13,6 +12,7 @@ import (
12 cxt "github.com/ipfs/go-ipfs/Godeps/_workspace/src/golang.org/x/net/context"
13
14 mdag "github.com/ipfs/go-ipfs/merkledag"
15 + ft "github.com/ipfs/go-ipfs/unixfs"
16 uio "github.com/ipfs/go-ipfs/unixfs/io"
17 upb "github.com/ipfs/go-ipfs/unixfs/pb"
18 )
@@ -21,7 +21,8 @@ import (
21 // TODO: does this need to be configurable?
22 var DefaultBufSize = 1048576
23
24 -func DagArchive(ctx cxt.Context, nd *mdag.Node, name string, dag mdag.DAGService, compression int) (io.Reader, error) {
24 +// DagArchive is equivalent to `ipfs getdag $hash | maybe_tar | maybe_gzip`
25 +func DagArchive(ctx cxt.Context, nd *mdag.Node, name string, dag mdag.DAGService, archive bool, compression int) (io.Reader, error) {
26
27 _, filename := path.Split(name)
28
@@ -31,17 +32,44 @@ func DagArchive(ctx cxt.Context, nd *mdag.Node, name string, dag mdag.DAGService
32 // use a buffered writer to parallelize task
33 bufw := bufio.NewWriterSize(pipew, DefaultBufSize)
34
35 + // compression determines whether to use gzip compression.
36 + var maybeGzw io.Writer
37 + if compression != gzip.NoCompression {
38 + var err error
39 + maybeGzw, err = gzip.NewWriterLevel(bufw, compression)
40 + if err != nil {
41 + return nil, err
42 + }
43 + } else {
44 + maybeGzw = bufw
45 + }
46 +
47 // construct the tar writer
35 - w, err := NewWriter(bufw, dag, compression)
48 + w, err := NewWriter(ctx, dag, archive, compression, maybeGzw)
49 if err != nil {
50 return nil, err
51 }
52
53 // write all the nodes recursively
54 go func() {
42 - if err := w.WriteNode(ctx, nd, filename); err != nil {
43 - pipew.CloseWithError(err)
44 - return
55 + if !archive && compression != gzip.NoCompression {
56 + // the case when the node is a file
57 + dagr, err := uio.NewDagReader(w.ctx, nd, w.Dag)
58 + if err != nil {
59 + pipew.CloseWithError(err)
60 + return
61 + }
62 +
63 + if _, err := dagr.WriteTo(maybeGzw); err != nil {
64 + pipew.CloseWithError(err)
65 + return
66 + }
67 + } else {
68 + // the case for 1. archive, and 2. not archived and not compressed, in which tar is used anyway as a transport format
69 + if err := w.WriteNode(nd, filename); err != nil {
70 + pipew.CloseWithError(err)
71 + return
72 + }
73 }
74
75 if err := bufw.Flush(); err != nil {
@@ -49,6 +77,7 @@ func DagArchive(ctx cxt.Context, nd *mdag.Node, name string, dag mdag.DAGService
77 return
78 }
79
80 + w.Close()
81 pipew.Close() // everything seems to be ok.
82 }()
83
@@ -61,39 +90,32 @@ func DagArchive(ctx cxt.Context, nd *mdag.Node, name string, dag mdag.DAGService
90 type Writer struct {
91 Dag mdag.DAGService
92 TarW *tar.Writer
93 +
94 + ctx cxt.Context
95 }
96
97 // NewWriter wraps given io.Writer.
67 -// compression determines whether to use gzip compression.
68 -func NewWriter(w io.Writer, dag mdag.DAGService, compression int) (*Writer, error) {
69 -
70 - if compression != gzip.NoCompression {
71 - var err error
72 - w, err = gzip.NewWriterLevel(w, compression)
73 - if err != nil {
74 - return nil, err
75 - }
76 - }
77 -
98 +func NewWriter(ctx cxt.Context, dag mdag.DAGService, archive bool, compression int, w io.Writer) (*Writer, error) {
99 return &Writer{
100 Dag: dag,
101 TarW: tar.NewWriter(w),
102 + ctx: ctx,
103 }, nil
104 }
105
84 -func (w *Writer) WriteDir(ctx cxt.Context, nd *mdag.Node, fpath string) error {
106 +func (w *Writer) writeDir(nd *mdag.Node, fpath string) error {
107 if err := writeDirHeader(w.TarW, fpath); err != nil {
108 return err
109 }
110
89 - for i, ng := range w.Dag.GetDAG(ctx, nd) {
90 - child, err := ng.Get(ctx)
111 + for i, ng := range w.Dag.GetDAG(w.ctx, nd) {
112 + child, err := ng.Get(w.ctx)
113 if err != nil {
114 return err
115 }
116
117 npath := path.Join(fpath, nd.Links[i].Name)
96 - if err := w.WriteNode(ctx, child, npath); err != nil {
118 + if err := w.WriteNode(child, npath); err != nil {
119 return err
120 }
121 }
@@ -101,46 +123,33 @@ func (w *Writer) WriteDir(ctx cxt.Context, nd *mdag.Node, fpath string) error {
123 return nil
124 }
125
104 -func (w *Writer) WriteFile(ctx cxt.Context, nd *mdag.Node, fpath string) error {
105 - pb := new(upb.Data)
106 - if err := proto.Unmarshal(nd.Data, pb); err != nil {
107 - return err
108 - }
109 -
110 - return w.writeFile(ctx, nd, pb, fpath)
111 -}
112 -
113 -func (w *Writer) writeFile(ctx cxt.Context, nd *mdag.Node, pb *upb.Data, fpath string) error {
126 +func (w *Writer) writeFile(nd *mdag.Node, pb *upb.Data, fpath string) error {
127 if err := writeFileHeader(w.TarW, fpath, pb.GetFilesize()); err != nil {
128 return err
129 }
130
118 - dagr, err := uio.NewDagReader(ctx, nd, w.Dag)
119 - if err != nil {
120 - return err
121 - }
122 -
123 - _, err = io.Copy(w.TarW, dagr)
124 - if err != nil && err != io.EOF {
125 - return err
126 - }
127 -
128 - return nil
131 + dagr := uio.NewDataFileReader(w.ctx, nd, pb, w.Dag)
132 + _, err := dagr.WriteTo(w.TarW)
133 + return err
134 }
135
131 -func (w *Writer) WriteNode(ctx cxt.Context, nd *mdag.Node, fpath string) error {
136 +func (w *Writer) WriteNode(nd *mdag.Node, fpath string) error {
137 pb := new(upb.Data)
138 if err := proto.Unmarshal(nd.Data, pb); err != nil {
139 return err
140 }
141
142 switch pb.GetType() {
143 + case upb.Data_Metadata:
144 + fallthrough
145 case upb.Data_Directory:
139 - return w.WriteDir(ctx, nd, fpath)
146 + return w.writeDir(nd, fpath)
147 + case upb.Data_Raw:
148 + fallthrough
149 case upb.Data_File:
141 - return w.writeFile(ctx, nd, pb, fpath)
150 + return w.writeFile(nd, pb, fpath)
151 default:
143 - return fmt.Errorf("unixfs type not supported: %s", pb.GetType())
152 + return ft.ErrUnrecognizedType
153 }
154 }
155