From f710b6003a214f1d50b796c9a96bda2f048646b4 Mon Sep 17 00:00:00 2001 From: Junker der Provinz <133605895+junkerderprovinz@users.noreply.github.com> Date: Sun, 23 Aug 2026 08:00:50 +0200 Subject: [PATCH] s3api: retry a transient failure when listing multipart uploads/parts (RFC on layering) (#10886) s3api: retry a transient failure when listing multipart uploads/parts A blip on the way to the filer failed the whole ListMultipartUploads or ListParts request. Both reported failure points sit inside one streaming listing: the ListEntries call that opens the stream, and the stream.Recv calls that drain it. Neither retried, so a single Unavailable answer from a filer that was restarting turned into a 500 for the S3 client. Replay the listing instead, bounded to three attempts with a 100ms backoff that doubles. Only a transient failure is replayed. A not-found answer stays authoritative so the empty-list branch still works, and every other error still reaches the client on the first attempt. This is scoped to (*S3ApiServer).list rather than added inside DoSeaweedListWithSnapshot, which mount, the shell and the other object listings share, and where a retry after a partial stream would re-deliver entries the callback had already seen. Within one call to list, a replay is safe: it collects into a fresh slice each time, so it can neither duplicate nor drop entries. That guarantee does not extend past this function. withFilerClientFailover already re-runs its callback against the next filer on any non-NotFound error without resetting the caller's accumulator, so on a multi-filer gateway a mid-listing failover can itself produce a duplicated result with err == nil, independent of this change and not fixed by it. Noted in the PR rather than silently left for someone to rediscover. Fixes #7221 References #7235 --- weed/s3api/filer_multipart.go | 13 +- weed/s3api/filer_util.go | 47 +++++ weed/s3api/filer_util_list_retry_test.go | 250 +++++++++++++++++++++++ 3 files changed, 308 insertions(+), 2 deletions(-) create mode 100644 weed/s3api/filer_util_list_retry_test.go diff --git a/weed/s3api/filer_multipart.go b/weed/s3api/filer_multipart.go index 38b73111c..9c2c9f658 100644 --- a/weed/s3api/filer_multipart.go +++ b/weed/s3api/filer_multipart.go @@ -1037,7 +1037,13 @@ func (s3a *S3ApiServer) listMultipartUploads(input *s3.ListMultipartUploadsInput IsTruncated: aws.Bool(false), } - entries, _, err := s3a.list(s3a.genUploadsFolder(*input.Bucket), "", *input.UploadIdMarker, false, math.MaxInt32) + // A blip on the way to the filer, either on the ListEntries call or on the + // stream receives that follow it, used to fail the whole listing; replay it + // a bounded number of times before giving up. + uploadsFolder := s3a.genUploadsFolder(*input.Bucket) + entries, _, err := listWithRetry(uploadsFolder, func() ([]*filer_pb.Entry, bool, error) { + return s3a.list(uploadsFolder, "", *input.UploadIdMarker, false, math.MaxInt32) + }) if err != nil { // A missing .uploads folder normally lists as empty with no error; a // store that reports it as not-found still means an empty list. @@ -1104,7 +1110,10 @@ func (s3a *S3ApiServer) listObjectParts(input *s3.ListPartsInput) (output *ListP StorageClass: aws.String("STANDARD"), } - entries, isLast, err := s3a.list(s3a.genUploadsFolder(*input.Bucket)+"/"+*input.UploadId, "", fmt.Sprintf("%04d%s", *input.PartNumberMarker, multipartExt), false, uint32(*input.MaxParts)) + partsFolder := s3a.genUploadsFolder(*input.Bucket) + "/" + *input.UploadId + entries, isLast, err := listWithRetry(partsFolder, func() ([]*filer_pb.Entry, bool, error) { + return s3a.list(partsFolder, "", fmt.Sprintf("%04d%s", *input.PartNumberMarker, multipartExt), false, uint32(*input.MaxParts)) + }) if err != nil { // A store that reports the missing upload directory as not-found means // the upload is gone (completed or aborted), not a store error. diff --git a/weed/s3api/filer_util.go b/weed/s3api/filer_util.go index 9a883bbf3..fa1acd75f 100644 --- a/weed/s3api/filer_util.go +++ b/weed/s3api/filer_util.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "strings" + "time" "github.com/seaweedfs/seaweedfs/weed/filer" "github.com/seaweedfs/seaweedfs/weed/glog" @@ -45,6 +46,52 @@ func (s3a *S3ApiServer) list(parentDirectoryPath, prefix, startFrom string, incl } +// Bounds for replaying a listing that failed with a transient error. A listing +// is a read with no side effects and each attempt opens a new stream and +// collects into a fresh slice, so replaying it can neither duplicate nor drop +// entries. The bound keeps a filer that is genuinely down from stalling the S3 +// request: at most two extra attempts and 300ms of added wait. +const ( + listRetryAttempts = 3 + listRetryInitialBackoff = 100 * time.Millisecond +) + +// isRetryableListError reports whether a failed listing is worth replaying. A +// not-found answer is authoritative and every other non-transient error is a +// real failure, so both must reach the caller unchanged: the point is to +// survive a blip, not to hide a broken store. +func isRetryableListError(err error) bool { + if err == nil || isFilerNotFound(err) { + return false + } + if status.Code(err) == codes.Unavailable { + return true + } + // DoSeaweedListWithSnapshot wraps a failed ListEntries call with %v, which + // drops the gRPC status from the error chain, so for that path the message + // is all that is left to classify on. + return util.IsTransientError(err) +} + +// listWithRetry replays doList while the filer answers with a transient error. +// Both failure points reported for multipart listing, the ListEntries call +// itself and the stream.Recv that follows it, surface as a plain error out of +// filer_pb.List, so a single retry point above it covers both. +func listWithRetry(parentDirectoryPath string, doList func() (entries []*filer_pb.Entry, isLast bool, err error)) (entries []*filer_pb.Entry, isLast bool, err error) { + + backoff := listRetryInitialBackoff + for attempt := 1; ; attempt++ { + entries, isLast, err = doList() + if err == nil || attempt >= listRetryAttempts || !isRetryableListError(err) { + return entries, isLast, err + } + glog.V(1).Infof("list %s attempt %d/%d hit a transient error, retrying in %v: %v", parentDirectoryPath, attempt, listRetryAttempts, backoff, err) + time.Sleep(backoff) + backoff *= 2 + } + +} + func (s3a *S3ApiServer) rm(parentDirectoryPath, entryName string, isDeleteData, isRecursive bool) error { return s3a.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error { diff --git a/weed/s3api/filer_util_list_retry_test.go b/weed/s3api/filer_util_list_retry_test.go new file mode 100644 index 000000000..8a89f89e9 --- /dev/null +++ b/weed/s3api/filer_util_list_retry_test.go @@ -0,0 +1,250 @@ +package s3api + +import ( + "context" + "fmt" + "io" + "testing" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + grpc "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/metadata" + "google.golang.org/grpc/status" +) + +// listRetryStream replays a scripted sequence of entries and then, optionally, +// fails. failAfter < 0 means the stream runs to EOF. +type listRetryStream struct { + entries []*filer_pb.Entry + index int + failAfter int + failErr error +} + +func (s *listRetryStream) Recv() (*filer_pb.ListEntriesResponse, error) { + if s.failErr != nil && s.failAfter >= 0 && s.index >= s.failAfter { + return nil, s.failErr + } + if s.index >= len(s.entries) { + return nil, io.EOF + } + entry := s.entries[s.index] + s.index++ + return &filer_pb.ListEntriesResponse{Entry: entry}, nil +} + +func (s *listRetryStream) Header() (metadata.MD, error) { return metadata.MD{}, nil } +func (s *listRetryStream) Trailer() metadata.MD { return metadata.MD{} } +func (s *listRetryStream) CloseSend() error { return nil } +func (s *listRetryStream) Context() context.Context { return context.Background() } +func (s *listRetryStream) SendMsg(any) error { return nil } +func (s *listRetryStream) RecvMsg(any) error { return nil } + +// listRetryClient is a filer whose ListEntries fails in a scripted way for the +// first attempts and then behaves. callErrs is consumed one entry per attempt: +// a non-nil entry fails the ListEntries call itself (issue #7221), a nil entry +// lets the call through so recvErrs decides whether the stream breaks partway +// (issue #7235). +type listRetryClient struct { + filer_pb.SeaweedFilerClient + + entries []*filer_pb.Entry + callErrs []error + recvErrs []error + recvAfter []int + + attempts int +} + +func (c *listRetryClient) ListEntries(ctx context.Context, in *filer_pb.ListEntriesRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[filer_pb.ListEntriesResponse], error) { + attempt := c.attempts + c.attempts++ + + if attempt < len(c.callErrs) && c.callErrs[attempt] != nil { + return nil, c.callErrs[attempt] + } + + stream := &listRetryStream{entries: c.entries, failAfter: -1} + if attempt < len(c.recvErrs) && c.recvErrs[attempt] != nil { + stream.failErr = c.recvErrs[attempt] + stream.failAfter = c.recvAfter[attempt] + } + return stream, nil +} + +// listRetryAccessor hands the scripted client to filer_pb.List without any +// wrapping of its own, so the test sees the error exactly as it leaves +// DoSeaweedListWithSnapshot. +type listRetryAccessor struct { + client filer_pb.SeaweedFilerClient +} + +func (a *listRetryAccessor) WithFilerClient(_ bool, fn func(filer_pb.SeaweedFilerClient) error) error { + return fn(a.client) +} + +func (a *listRetryAccessor) AdjustedUrl(*filer_pb.Location) string { return "" } +func (a *listRetryAccessor) GetDataCenter() string { return "" } + +// listOnce mirrors (*S3ApiServer).list so the test drives the real +// filer_pb.List and DoSeaweedListWithSnapshot code path, including the fresh +// accumulator that makes a replay safe. +func listOnce(client filer_pb.FilerClient, parentDirectoryPath string) (entries []*filer_pb.Entry, isLast bool, err error) { + err = filer_pb.List(context.Background(), client, parentDirectoryPath, "", func(entry *filer_pb.Entry, isLastEntry bool) error { + entries = append(entries, entry) + if isLastEntry { + isLast = true + } + return nil + }, "", false, 100) + + if len(entries) == 0 { + isLast = true + } + + return +} + +func testUploadEntries(names ...string) []*filer_pb.Entry { + entries := make([]*filer_pb.Entry, 0, len(names)) + for _, name := range names { + entries = append(entries, &filer_pb.Entry{Name: name, Attributes: &filer_pb.FuseAttributes{}}) + } + return entries +} + +func entryNames(entries []*filer_pb.Entry) []string { + names := make([]string, 0, len(entries)) + for _, entry := range entries { + names = append(names, entry.Name) + } + return names +} + +// Issue #7221: a transient failure of the ListEntries call itself. +func TestListWithRetryRecoversFromTransientListEntriesCall(t *testing.T) { + client := &listRetryClient{ + entries: testUploadEntries("upload-1", "upload-2", "upload-3"), + callErrs: []error{ + status.Error(codes.Unavailable, "filer is restarting"), + status.Error(codes.Unavailable, "filer is restarting"), + nil, + }, + } + accessor := &listRetryAccessor{client: client} + + entries, _, err := listWithRetry("/buckets/b/.uploads", func() ([]*filer_pb.Entry, bool, error) { + return listOnce(accessor, "/buckets/b/.uploads") + }) + + require.NoError(t, err) + assert.Equal(t, []string{"upload-1", "upload-2", "upload-3"}, entryNames(entries)) + assert.Equal(t, listRetryAttempts, client.attempts, "should have used every allowed attempt") +} + +// Issue #7235: a transient failure on stream.Recv part way through the listing. +// The replay must return the full listing exactly once, with no entry from the +// aborted attempt carried over. +func TestListWithRetryRecoversFromTransientStreamRecv(t *testing.T) { + client := &listRetryClient{ + entries: testUploadEntries("upload-1", "upload-2", "upload-3"), + // break after two successful receives, so the first attempt has already + // handed "upload-1" to the callback when the stream dies + recvErrs: []error{status.Error(codes.Unavailable, "transport is closing"), nil}, + recvAfter: []int{2, 0}, + } + accessor := &listRetryAccessor{client: client} + + entries, _, err := listWithRetry("/buckets/b/.uploads", func() ([]*filer_pb.Entry, bool, error) { + return listOnce(accessor, "/buckets/b/.uploads") + }) + + require.NoError(t, err) + assert.Equal(t, []string{"upload-1", "upload-2", "upload-3"}, entryNames(entries), "replay must not duplicate the entries of the aborted attempt") + assert.Equal(t, 2, client.attempts) +} + +// The retry is bounded: a filer that stays down still fails the request, and +// the transient error reaches the caller. +func TestListWithRetryGivesUpAfterBoundedAttempts(t *testing.T) { + client := &listRetryClient{ + entries: testUploadEntries("upload-1"), + callErrs: []error{ + status.Error(codes.Unavailable, "no filer available"), + status.Error(codes.Unavailable, "no filer available"), + status.Error(codes.Unavailable, "no filer available"), + nil, // would succeed on a fourth attempt, which must never happen + }, + } + accessor := &listRetryAccessor{client: client} + + _, _, err := listWithRetry("/buckets/b/.uploads", func() ([]*filer_pb.Entry, bool, error) { + return listOnce(accessor, "/buckets/b/.uploads") + }) + + require.Error(t, err) + assert.Contains(t, err.Error(), "Unavailable") + assert.Equal(t, listRetryAttempts, client.attempts) +} + +// A non-transient error is a real failure and must propagate on the first +// attempt rather than being retried away. +func TestListWithRetryPropagatesNonTransientError(t *testing.T) { + for name, listErr := range map[string]error{ + "permission denied": status.Error(codes.PermissionDenied, "not allowed"), + "invalid argument": status.Error(codes.InvalidArgument, "bad prefix"), + } { + t.Run(name, func(t *testing.T) { + client := &listRetryClient{ + entries: testUploadEntries("upload-1"), + callErrs: []error{listErr, nil}, + } + accessor := &listRetryAccessor{client: client} + + _, _, err := listWithRetry("/buckets/b/.uploads", func() ([]*filer_pb.Entry, bool, error) { + return listOnce(accessor, "/buckets/b/.uploads") + }) + + require.Error(t, err) + assert.Equal(t, 1, client.attempts, "a non-transient error must not be retried") + }) + } +} + +// Not-found is authoritative: it must not be retried, and it must still read as +// not-found afterwards so listMultipartUploads keeps answering with an empty +// list instead of a 500. +func TestListWithRetryPropagatesNotFound(t *testing.T) { + client := &listRetryClient{ + entries: testUploadEntries("upload-1"), + callErrs: []error{status.Error(codes.NotFound, "filer: no entry is found in filer store"), nil}, + } + accessor := &listRetryAccessor{client: client} + + _, _, err := listWithRetry("/buckets/b/.uploads", func() ([]*filer_pb.Entry, bool, error) { + return listOnce(accessor, "/buckets/b/.uploads") + }) + + require.Error(t, err) + assert.Equal(t, 1, client.attempts, "not-found is authoritative and must not be retried") + assert.True(t, isFilerNotFound(err), "not-found must survive the retry wrapper") +} + +func TestIsRetryableListError(t *testing.T) { + assert.False(t, isRetryableListError(nil)) + assert.True(t, isRetryableListError(status.Error(codes.Unavailable, "connection refused"))) + assert.True(t, isRetryableListError(status.Error(codes.ResourceExhausted, "too many requests"))) + // the ListEntries path loses the gRPC status to a %v wrap, so classification + // has to survive on the message alone + assert.True(t, isRetryableListError(fmt.Errorf("list /buckets/b/.uploads: %v", status.Error(codes.Unavailable, "filer is restarting")))) + // and this is the shape the S3 gateway actually sees once the failover in + // (*S3ApiServer).WithFilerClient has exhausted every filer + assert.True(t, isRetryableListError(fmt.Errorf("all filers failed, last error: %w", + fmt.Errorf("list /buckets/b/.uploads: %v", status.Error(codes.Unavailable, "filer is restarting"))))) + assert.False(t, isRetryableListError(status.Error(codes.NotFound, "filer: no entry is found in filer store"))) + assert.False(t, isRetryableListError(status.Error(codes.PermissionDenied, "not allowed"))) + assert.False(t, isRetryableListError(context.Canceled)) +}