@cryptotaxi247 / kubo / commits / ccb693ab4

teach unixfs/tar (aka ipfs get) how to use GetBlocks

Jeromy committed Feb 18, 2015 at 08:44 UTC ccb693ab4d53553ac095998238a465fa172f133b
1 file changed +49 -47
unixfs/tar/reader.go
+49 -47
@@ -60,20 +60,20 @@ func NewReader(path path.Path, dag mdag.DAGService, resolver *path.Resolver, com
60 return reader, nil
61 }
62
63 -func (i *Reader) writeToBuf(dagnode *mdag.Node, path string, depth int) {
63 +func (r *Reader) writeToBuf(dagnode *mdag.Node, path string, depth int) {
64 pb := new(upb.Data)
65 err := proto.Unmarshal(dagnode.Data, pb)
66 if err != nil {
67 - i.emitError(err)
67 + r.emitError(err)
68 return
69 }
70
71 if depth == 0 {
72 - defer i.close()
72 + defer r.close()
73 }
74
75 if pb.GetType() == upb.Data_Directory {
76 - err = i.writer.WriteHeader(&tar.Header{
76 + err = r.writer.WriteHeader(&tar.Header{
77 Name: path,
78 Typeflag: tar.TypeDir,
79 Mode: 0777,
@@ -81,23 +81,25 @@ func (i *Reader) writeToBuf(dagnode *mdag.Node, path string, depth int) {
81 // TODO: set mode, dates, etc. when added to unixFS
82 })
83 if err != nil {
84 - i.emitError(err)
84 + r.emitError(err)
85 return
86 }
87 - i.flush()
87 + r.flush()
88
89 - for _, link := range dagnode.Links {
90 - childNode, err := link.GetNode(i.dag)
89 + ctx, _ := context.WithTimeout(context.TODO(), time.Second*60)
90 +
91 + for i, ng := range r.dag.GetDAG(ctx, dagnode) {
92 + childNode, err := ng.Get()
93 if err != nil {
92 - i.emitError(err)
94 + r.emitError(err)
95 return
96 }
95 - i.writeToBuf(childNode, gopath.Join(path, link.Name), depth+1)
97 + r.writeToBuf(childNode, gopath.Join(path, dagnode.Links[i].Name), depth+1)
98 }
99 return
100 }
101
100 - err = i.writer.WriteHeader(&tar.Header{
102 + err = r.writer.WriteHeader(&tar.Header{
103 Name: path,
104 Size: int64(pb.GetFilesize()),
105 Typeflag: tar.TypeReg,
@@ -106,95 +108,95 @@ func (i *Reader) writeToBuf(dagnode *mdag.Node, path string, depth int) {
108 // TODO: set mode, dates, etc. when added to unixFS
109 })
110 if err != nil {
109 - i.emitError(err)
111 + r.emitError(err)
112 return
113 }
112 - i.flush()
114 + r.flush()
115
114 - reader, err := uio.NewDagReader(context.TODO(), dagnode, i.dag)
116 + reader, err := uio.NewDagReader(context.TODO(), dagnode, r.dag)
117 if err != nil {
116 - i.emitError(err)
118 + r.emitError(err)
119 return
120 }
121
120 - err = i.syncCopy(reader)
122 + err = r.syncCopy(reader)
123 if err != nil {
122 - i.emitError(err)
124 + r.emitError(err)
125 return
126 }
127 }
128
127 -func (i *Reader) Read(p []byte) (int, error) {
129 +func (r *Reader) Read(p []byte) (int, error) {
130 // wait for the goroutine that is writing data to the buffer to tell us
131 // there is something to read
130 - if !i.closed {
131 - <-i.signalChan
132 + if !r.closed {
133 + <-r.signalChan
134 }
135
134 - if i.err != nil {
135 - return 0, i.err
136 + if r.err != nil {
137 + return 0, r.err
138 }
139
138 - if !i.closed {
139 - defer i.signal()
140 + if !r.closed {
141 + defer r.signal()
142 }
143
142 - if i.buf.Len() == 0 {
143 - if i.closed {
144 + if r.buf.Len() == 0 {
145 + if r.closed {
146 return 0, io.EOF
147 }
148 return 0, nil
149 }
150
149 - n, err := i.buf.Read(p)
150 - if err == io.EOF && !i.closed || i.buf.Len() > 0 {
151 + n, err := r.buf.Read(p)
152 + if err == io.EOF && !r.closed || r.buf.Len() > 0 {
153 return n, nil
154 }
155
156 return n, err
157 }
158
157 -func (i *Reader) signal() {
158 - i.signalChan <- struct{}{}
159 +func (r *Reader) signal() {
160 + r.signalChan <- struct{}{}
161 }
162
161 -func (i *Reader) flush() {
162 - i.signal()
163 - <-i.signalChan
163 +func (r *Reader) flush() {
164 + r.signal()
165 + <-r.signalChan
166 }
167
166 -func (i *Reader) emitError(err error) {
167 - i.err = err
168 - i.signal()
168 +func (r *Reader) emitError(err error) {
169 + r.err = err
170 + r.signal()
171 }
172
171 -func (i *Reader) close() {
172 - i.closed = true
173 - defer i.signal()
174 - err := i.writer.Close()
173 +func (r *Reader) close() {
174 + r.closed = true
175 + defer r.signal()
176 + err := r.writer.Close()
177 if err != nil {
176 - i.emitError(err)
178 + r.emitError(err)
179 return
180 }
179 - if i.gzipWriter != nil {
180 - err = i.gzipWriter.Close()
181 + if r.gzipWriter != nil {
182 + err = r.gzipWriter.Close()
183 if err != nil {
182 - i.emitError(err)
184 + r.emitError(err)
185 return
186 }
187 }
188 }
189
188 -func (i *Reader) syncCopy(reader io.Reader) error {
190 +func (r *Reader) syncCopy(reader io.Reader) error {
191 buf := make([]byte, 32*1024)
192 for {
193 nr, err := reader.Read(buf)
194 if nr > 0 {
193 - _, err := i.writer.Write(buf[:nr])
195 + _, err := r.writer.Write(buf[:nr])
196 if err != nil {
197 return err
198 }
197 - i.flush()
199 + r.flush()
200 }
201 if err == io.EOF {
202 break