Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions storage/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1000,13 +1000,23 @@ func TestOpenAppendableWriterMultipleFlushesEmulated(t *testing.T) {
w.Append = true
// This should chunk the request into three separate flushes to storage.
w.ChunkSize = MiB
var lastReportedOffset int64
w.ProgressFunc = func(offset int64) {
if offset != lastReportedOffset+MiB {
t.Errorf("incorrect progress report: got %d; want %d", offset, lastReportedOffset+MiB)
}
lastReportedOffset = offset
}
_, err = w.Write(randomBytes3MiB)
if err != nil {
t.Fatalf("writing test data: got %v; want ok", err)
}
if err := w.Close(); err != nil {
t.Fatalf("closing test data writer: got %v; want ok", err)
}
if lastReportedOffset != 3*MiB {
t.Errorf("incorrect final progress report: got %d; want %d", lastReportedOffset, 3*MiB)
}

if diff := cmp.Diff(w.Attrs().Name, objName); diff != "" {
t.Fatalf("Resulting object name: got(-), want(+):\n%s", diff)
Expand Down
21 changes: 16 additions & 5 deletions storage/grpc_writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ type gRPCAppendBidiWriteBufferSender struct {
objectChecksums *storagepb.ObjectChecksums

forceFirstMessage bool
progress func(int64)
flushOffset int64

// Fields used to report responses from the receive side of the stream
Expand All @@ -62,6 +63,7 @@ func (w *gRPCWriter) newGRPCAppendBidiWriteBufferSender() (*gRPCAppendBidiWriteB
},
objectChecksums: toProtoChecksums(w.sendCRC32C, w.attrs),
forceFirstMessage: true,
progress: w.progress,
}
return s, nil
}
Expand Down Expand Up @@ -246,26 +248,35 @@ func (s *gRPCAppendBidiWriteBufferSender) sendOnConnectedStream(buf []byte, offs
if s.recvErr != io.EOF {
return nil, s.recvErr
}
if obj.GetSize() > s.flushOffset {
s.flushOffset = obj.GetSize()
s.progress(s.flushOffset)
}
return
}

if flush {
// We don't necessarily expect multiple responses for a single flush, but
// this allows the server to send multiple responses if it wants to.
for s.flushOffset < offset+int64(len(buf)) {
flushOffset := s.flushOffset
for flushOffset < offset+int64(len(buf)) {
resp, ok := <-s.recvs
if !ok {
return nil, s.recvErr
}
pSize := resp.GetPersistedSize()
rSize := resp.GetResource().GetSize()
if s.flushOffset < pSize {
s.flushOffset = pSize
if flushOffset < pSize {
flushOffset = pSize
}
if s.flushOffset < rSize {
s.flushOffset = rSize
if flushOffset < rSize {
flushOffset = rSize
}
}
if s.flushOffset < flushOffset {
s.flushOffset = flushOffset
s.progress(s.flushOffset)
}
}

return
Expand Down