core/commands: Made 'get' copying thread-safe
Matt Bell committed
Jan 22, 2015 at 10:05 UTC
276a8d062cd76fc8875af2fd849312a584e7e187
1 file changed
+108
-67
core/commands/get.go
+108
-67
@@ -5,6 +5,7 @@ import (
5
"bytes"
6
"io"
7
p "path"
8
+ "sync"
9
10
cmds "github.com/jbenet/go-ipfs/commands"
11
core "github.com/jbenet/go-ipfs/core"
@@ -42,137 +43,177 @@ To output a TAR archive instead of unpacked files, use '--archive' or '-a'.
43
return
44
}
45
45
- reader, err := get(node, req.Arguments())
46
+ reader, err := get(node, req.Arguments()[0])
47
if err != nil {
48
res.SetError(err, cmds.ErrNormal)
49
return
50
}
51
res.SetOutput(reader)
52
},
52
-
53
- // TODO: create a PostRun that splits the archive up into files
53
}
54
56
-func get(node *core.IpfsNode, paths []string) (io.Reader, error) {
57
- reader := &getReader{signalChan: make(chan struct{})}
58
- writer := tar.NewWriter(&reader.buf)
55
+func get(node *core.IpfsNode, path string) (io.Reader, error) {
56
+ buf := NewBufReadWriter()
57
58
go func() {
61
- for _, path := range paths {
62
- _, err := copyFile(node, writer, path, nil, reader.signalChan)
63
- if err != nil {
64
- log.Error(err)
65
- return
66
- }
67
- }
68
-
69
- err := writer.Flush()
59
+ err := copyFilesAsTar(node, buf, path)
60
if err != nil {
61
log.Error(err)
62
return
63
}
74
-
75
- reader.Close()
76
- reader.Signal()
64
}()
65
79
- return reader, nil
66
+ return buf, nil
67
}
68
82
-func copyFile(node *core.IpfsNode, writer *tar.Writer, path string, dagnode *dag.Node, signal chan struct{}) (int64, error) {
69
+func copyFilesAsTar(node *core.IpfsNode, buf *bufReadWriter, path string) error {
70
+ writer := tar.NewWriter(buf)
71
+
72
+ err := _copyFilesAsTar(node, writer, buf, path, nil)
73
+ if err != nil {
74
+ return err
75
+ }
76
+
77
+ err = writer.Flush()
78
+ if err != nil {
79
+ return err
80
+ }
81
+ buf.Close()
82
+ buf.Signal()
83
+ return nil
84
+}
85
+
86
+func _copyFilesAsTar(node *core.IpfsNode, writer *tar.Writer, buf *bufReadWriter, path string, dagnode *dag.Node) error {
87
var err error
88
if dagnode == nil {
89
dagnode, err = node.Resolver.ResolvePath(path)
90
if err != nil {
87
- return 0, err
91
+ return err
92
}
93
}
94
95
pb := new(upb.Data)
96
err = proto.Unmarshal(dagnode.Data, pb)
97
if err != nil {
94
- return 0, err
98
+ return err
99
}
100
97
- written := int64(0)
101
if pb.GetType() == upb.Data_Directory {
102
+ buf.mutex.Lock()
103
err = writer.WriteHeader(&tar.Header{
104
Name: path,
105
Typeflag: tar.TypeDir,
106
Mode: 0777,
107
// TODO: set mode, dates, etc. when added to unixFS
108
})
109
+ buf.mutex.Unlock()
110
if err != nil {
106
- return 0, err
111
+ return err
112
}
113
114
for _, link := range dagnode.Links {
110
- n, err := copyFile(node, writer, p.Join(path, link.Name), link.Node, signal)
115
+ err := _copyFilesAsTar(node, writer, buf, p.Join(path, link.Name), link.Node)
116
if err != nil {
112
- return 0, err
117
+ return err
118
}
114
- written += n
119
}
116
- return written, nil
120
118
- } else {
119
- err = writer.WriteHeader(&tar.Header{
120
- Name: path,
121
- Size: int64(pb.GetFilesize()),
122
- Typeflag: tar.TypeReg,
123
- Mode: 0644,
124
- // TODO: set mode, dates, etc. when added to unixFS
125
- })
126
- if err != nil {
127
- return 0, err
128
- }
121
+ return nil
122
+ }
123
130
- reader, err := uio.NewDagReader(dagnode, node.DAG)
131
- if err != nil {
132
- return 0, err
133
- }
124
+ buf.mutex.Lock()
125
+ err = writer.WriteHeader(&tar.Header{
126
+ Name: path,
127
+ Size: int64(pb.GetFilesize()),
128
+ Typeflag: tar.TypeReg,
129
+ Mode: 0644,
130
+ // TODO: set mode, dates, etc. when added to unixFS
131
+ })
132
+ buf.mutex.Unlock()
133
+ if err != nil {
134
+ return err
135
+ }
136
135
- buf := make([]byte, 32*1024)
136
- for {
137
- nr, err := reader.Read(buf)
138
- if nr > 0 {
139
- nw, err := writer.Write(buf[:nr])
140
- if err != nil {
141
- return written, err
142
- }
143
- written += int64(nw)
144
- signal <- struct{}{}
145
- }
146
- if err == io.EOF {
147
- break
148
- }
149
- if err != nil {
150
- return written, err
151
- }
152
- }
153
- return written, nil
137
+ reader, err := uio.NewDagReader(dagnode, node.DAG)
138
+ if err != nil {
139
+ return err
140
+ }
141
+
142
+ _, err = syncCopy(writer, reader, buf)
143
+ if err != nil {
144
+ return err
145
}
146
+
147
+ return nil
148
}
149
157
-type getReader struct {
150
+type bufReadWriter struct {
151
buf bytes.Buffer
152
closed bool
153
signalChan chan struct{}
154
+ mutex *sync.Mutex
155
+}
156
+
157
+func NewBufReadWriter() *bufReadWriter {
158
+ return &bufReadWriter{
159
+ signalChan: make(chan struct{}),
160
+ mutex: &sync.Mutex{},
161
+ }
162
}
163
163
-func (i *getReader) Read(p []byte) (int, error) {
164
+func (i *bufReadWriter) Read(p []byte) (int, error) {
165
<-i.signalChan
166
+ i.mutex.Lock()
167
+ defer i.mutex.Unlock()
168
+
169
+ if i.buf.Len() == 0 {
170
+ if i.closed {
171
+ return 0, io.EOF
172
+ }
173
+ return 0, nil
174
+ }
175
+
176
n, err := i.buf.Read(p)
166
- if err == io.EOF && !i.closed {
177
+ if err == io.EOF && !i.closed || i.buf.Len() > 0 {
178
return n, nil
179
}
180
return n, err
181
}
182
172
-func (i *getReader) Signal() {
183
+func (i *bufReadWriter) Write(p []byte) (int, error) {
184
+ return i.buf.Write(p)
185
+}
186
+
187
+func (i *bufReadWriter) Signal() {
188
i.signalChan <- struct{}{}
189
}
190
176
-func (i *getReader) Close() {
191
+func (i *bufReadWriter) Close() error {
192
i.closed = true
193
+ return nil
194
+}
195
+
196
+func syncCopy(writer io.Writer, reader io.Reader, buf *bufReadWriter) (int64, error) {
197
+ written := int64(0)
198
+ copyBuf := make([]byte, 32*1024)
199
+ for {
200
+ nr, err := reader.Read(copyBuf)
201
+ if nr > 0 {
202
+ buf.mutex.Lock()
203
+ nw, err := writer.Write(copyBuf[:nr])
204
+ buf.mutex.Unlock()
205
+ if err != nil {
206
+ return written, err
207
+ }
208
+ written += int64(nw)
209
+ buf.Signal()
210
+ }
211
+ if err == io.EOF {
212
+ break
213
+ }
214
+ if err != nil {
215
+ return written, err
216
+ }
217
+ }
218
+ return written, nil
219
}