From 5095a2d6c9eb70831bcfb5f114052d6d721b67d0 Mon Sep 17 00:00:00 2001 From: Chris Carlon Date: Fri, 24 Jan 2025 16:05:32 +0000 Subject: [PATCH] fix(storage/append): Report progress for appends. The appendable writer should report progress, as the resumable writer does. --- storage/client_test.go | 10 ++++++++++ storage/grpc_writer.go | 21 ++++++++++++++++----- 2 files changed, 26 insertions(+), 5 deletions(-) diff --git a/storage/client_test.go b/storage/client_test.go index 88fa9d7c6567..7a45e0cc68a0 100644 --- a/storage/client_test.go +++ b/storage/client_test.go @@ -1000,6 +1000,13 @@ 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) @@ -1007,6 +1014,9 @@ func TestOpenAppendableWriterMultipleFlushesEmulated(t *testing.T) { 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) diff --git a/storage/grpc_writer.go b/storage/grpc_writer.go index 9c8e8fc30e6a..0d96cbe39b61 100644 --- a/storage/grpc_writer.go +++ b/storage/grpc_writer.go @@ -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 @@ -62,6 +63,7 @@ func (w *gRPCWriter) newGRPCAppendBidiWriteBufferSender() (*gRPCAppendBidiWriteB }, objectChecksums: toProtoChecksums(w.sendCRC32C, w.attrs), forceFirstMessage: true, + progress: w.progress, } return s, nil } @@ -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