Skip to content

Commit 618e300

Browse files
re-implement ReaderAt implementation with race protection (#1673)
Signed-off-by: bwplotka <[email protected]> Co-authored-by: Harshavardhana <[email protected]>
1 parent 4aec031 commit 618e300

File tree

2 files changed

+52
-46
lines changed

2 files changed

+52
-46
lines changed

api-put-object-streaming.go

+20-30
Original file line numberDiff line numberDiff line change
@@ -130,32 +130,32 @@ func (c *Client) putObjectMultipartStreamFromReadAt(ctx context.Context, bucketN
130130
var complMultipartUpload completeMultipartUpload
131131

132132
// Declare a channel that sends the next part number to be uploaded.
133-
// Buffered to 10000 because thats the maximum number of parts allowed
134-
// by S3.
135-
uploadPartsCh := make(chan uploadPartReq, 10000)
133+
uploadPartsCh := make(chan uploadPartReq)
136134

137135
// Declare a channel that sends back the response of a part upload.
138-
// Buffered to 10000 because thats the maximum number of parts allowed
139-
// by S3.
140-
uploadedPartsCh := make(chan uploadedPartRes, 10000)
136+
uploadedPartsCh := make(chan uploadedPartRes)
141137

142138
// Used for readability, lastPartNumber is always totalPartsCount.
143139
lastPartNumber := totalPartsCount
144140

141+
partitionCtx, partitionCancel := context.WithCancel(ctx)
142+
defer partitionCancel()
145143
// Send each part number to the channel to be processed.
146-
for p := 1; p <= totalPartsCount; p++ {
147-
uploadPartsCh <- uploadPartReq{PartNum: p}
148-
}
149-
close(uploadPartsCh)
150-
151-
partsBuf := make([][]byte, opts.getNumThreads())
152-
for i := range partsBuf {
153-
partsBuf[i] = make([]byte, 0, partSize)
154-
}
144+
go func() {
145+
defer close(uploadPartsCh)
146+
147+
for p := 1; p <= totalPartsCount; p++ {
148+
select {
149+
case <-partitionCtx.Done():
150+
return
151+
case uploadPartsCh <- uploadPartReq{PartNum: p}:
152+
}
153+
}
154+
}()
155155

156156
// Receive each part number from the channel allowing three parallel uploads.
157157
for w := 1; w <= opts.getNumThreads(); w++ {
158-
go func(w int, partSize int64) {
158+
go func(partSize int64) {
159159
for {
160160
var uploadReq uploadPartReq
161161
var ok bool
@@ -181,21 +181,11 @@ func (c *Client) putObjectMultipartStreamFromReadAt(ctx context.Context, bucketN
181181
partSize = lastPartSize
182182
}
183183

184-
n, rerr := readFull(io.NewSectionReader(reader, readOffset, partSize), partsBuf[w-1][:partSize])
185-
if rerr != nil && rerr != io.ErrUnexpectedEOF && rerr != io.EOF {
186-
uploadedPartsCh <- uploadedPartRes{
187-
Error: rerr,
188-
}
189-
// Exit the goroutine.
190-
return
191-
}
192-
193-
// Get a section reader on a particular offset.
194-
hookReader := newHook(bytes.NewReader(partsBuf[w-1][:n]), opts.Progress)
184+
sectionReader := newHook(io.NewSectionReader(reader, readOffset, partSize), opts.Progress)
195185

196186
// Proceed to upload the part.
197187
objPart, err := c.uploadPart(ctx, bucketName, objectName,
198-
uploadID, hookReader, uploadReq.PartNum,
188+
uploadID, sectionReader, uploadReq.PartNum,
199189
"", "", partSize,
200190
opts.ServerSideEncryption,
201191
!opts.DisableContentSha256,
@@ -218,7 +208,7 @@ func (c *Client) putObjectMultipartStreamFromReadAt(ctx context.Context, bucketN
218208
Part: uploadReq.Part,
219209
}
220210
}
221-
}(w, partSize)
211+
}(partSize)
222212
}
223213

224214
// Gather the responses as they occur and update any
@@ -229,12 +219,12 @@ func (c *Client) putObjectMultipartStreamFromReadAt(ctx context.Context, bucketN
229219
return UploadInfo{}, ctx.Err()
230220
case uploadRes := <-uploadedPartsCh:
231221
if uploadRes.Error != nil {
222+
232223
return UploadInfo{}, uploadRes.Error
233224
}
234225

235226
// Update the totalUploadedSize.
236227
totalUploadedSize += uploadRes.Size
237-
// Store the parts to be completed in order.
238228
complMultipartUpload.Parts = append(complMultipartUpload.Parts, CompletePart{
239229
ETag: uploadRes.Part.ETag,
240230
PartNumber: uploadRes.Part.PartNumber,

hook-reader.go

+32-16
Original file line numberDiff line numberDiff line change
@@ -20,20 +20,25 @@ package minio
2020
import (
2121
"fmt"
2222
"io"
23+
"sync"
2324
)
2425

2526
// hookReader hooks additional reader in the source stream. It is
2627
// useful for making progress bars. Second reader is appropriately
2728
// notified about the exact number of bytes read from the primary
2829
// source on each Read operation.
2930
type hookReader struct {
31+
mu sync.RWMutex
3032
source io.Reader
3133
hook io.Reader
3234
}
3335

3436
// Seek implements io.Seeker. Seeks source first, and if necessary
3537
// seeks hook if Seek method is appropriately found.
3638
func (hr *hookReader) Seek(offset int64, whence int) (n int64, err error) {
39+
hr.mu.Lock()
40+
defer hr.mu.Unlock()
41+
3742
// Verify for source has embedded Seeker, use it.
3843
sourceSeeker, ok := hr.source.(io.Seeker)
3944
if ok {
@@ -43,33 +48,41 @@ func (hr *hookReader) Seek(offset int64, whence int) (n int64, err error) {
4348
}
4449
}
4550

46-
// Verify if hook has embedded Seeker, use it.
47-
hookSeeker, ok := hr.hook.(io.Seeker)
48-
if ok {
49-
var m int64
50-
m, err = hookSeeker.Seek(offset, whence)
51-
if err != nil {
52-
return 0, err
53-
}
54-
if n != m {
55-
return 0, fmt.Errorf("hook seeker seeked %d bytes, expected source %d bytes", m, n)
51+
if hr.hook != nil {
52+
// Verify if hook has embedded Seeker, use it.
53+
hookSeeker, ok := hr.hook.(io.Seeker)
54+
if ok {
55+
var m int64
56+
m, err = hookSeeker.Seek(offset, whence)
57+
if err != nil {
58+
return 0, err
59+
}
60+
if n != m {
61+
return 0, fmt.Errorf("hook seeker seeked %d bytes, expected source %d bytes", m, n)
62+
}
5663
}
5764
}
65+
5866
return n, nil
5967
}
6068

6169
// Read implements io.Reader. Always reads from the source, the return
6270
// value 'n' number of bytes are reported through the hook. Returns
6371
// error for all non io.EOF conditions.
6472
func (hr *hookReader) Read(b []byte) (n int, err error) {
73+
hr.mu.RLock()
74+
defer hr.mu.RUnlock()
75+
6576
n, err = hr.source.Read(b)
6677
if err != nil && err != io.EOF {
6778
return n, err
6879
}
69-
// Progress the hook with the total read bytes from the source.
70-
if _, herr := hr.hook.Read(b[:n]); herr != nil {
71-
if herr != io.EOF {
72-
return n, herr
80+
if hr.hook != nil {
81+
// Progress the hook with the total read bytes from the source.
82+
if _, herr := hr.hook.Read(b[:n]); herr != nil {
83+
if herr != io.EOF {
84+
return n, herr
85+
}
7386
}
7487
}
7588
return n, err
@@ -79,7 +92,10 @@ func (hr *hookReader) Read(b []byte) (n int, err error) {
7992
// reports the data read from the source to the hook.
8093
func newHook(source, hook io.Reader) io.Reader {
8194
if hook == nil {
82-
return source
95+
return &hookReader{source: source}
96+
}
97+
return &hookReader{
98+
source: source,
99+
hook: hook,
83100
}
84-
return &hookReader{source, hook}
85101
}

0 commit comments

Comments
 (0)