Skip to content

Commit ac5e251

Browse files
authored
fix(storage): stop BucketExist scan after first result (#1031)
Signed-off-by: Xuewei Wang <xueweiwang.sallery@gmail.com>
1 parent de43489 commit ac5e251

2 files changed

Lines changed: 172 additions & 18 deletions

File tree

internal/storage/minio.go

Lines changed: 21 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -322,28 +322,31 @@ func (m *MinioClient) ListPrefix(ctx context.Context, prefix string, recursive b
322322

323323
// BucketExist checks if the bucket exists by listing a single object.
324324
// We use ListObjects instead of BucketExists (HEAD bucket) to minimize
325-
// required S3 permissions. MaxKeys=1 ensures only one object is fetched
326-
// to avoid iterating the entire bucket, which can timeout on large buckets.
325+
// required S3 permissions. MaxKeys=1 avoids iterating the entire bucket,
326+
// which can timeout on large buckets, while preserving caller context
327+
// propagation.
327328
func (m *MinioClient) BucketExist(ctx context.Context, prefix string) (bool, error) {
328-
opts := minio.ListObjectsOptions{Prefix: prefix, MaxKeys: 1}
329-
// MaxKeys=1 ensures at most one result, so the channel closes naturally.
330-
// We must drain the channel to avoid goroutine leaks (see minio-go docs).
331-
exists := true
332-
var retErr error
333-
for obj := range m.cli.ListObjects(ctx, m.cfg.Bucket, opts) {
334-
if obj.Err != nil {
335-
if minio.ToErrorResponse(obj.Err).Code == "NoSuchBucket" {
336-
exists = false
337-
} else {
338-
retErr = fmt.Errorf("storage: %s list objects %w", m.cfg.Provider, obj.Err)
339-
}
340-
}
329+
subCtx, cancel := context.WithCancel(ctx)
330+
defer cancel()
331+
332+
objCh := m.cli.ListObjects(subCtx, m.cfg.Bucket, minio.ListObjectsOptions{
333+
Prefix: prefix,
334+
MaxKeys: 1,
335+
})
336+
337+
obj, ok := <-objCh
338+
if !ok {
339+
return true, nil
341340
}
342341

343-
if retErr != nil {
344-
return false, retErr
342+
if obj.Err != nil {
343+
if minio.ToErrorResponse(obj.Err).Code == "NoSuchBucket" {
344+
return false, nil
345+
}
346+
return false, fmt.Errorf("storage: %s list objects %w", m.cfg.Provider, obj.Err)
345347
}
346-
return exists, nil
348+
349+
return true, nil
347350
}
348351

349352
func (m *MinioClient) CreateBucket(ctx context.Context) error {

internal/storage/minio_test.go

Lines changed: 151 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,19 @@
11
package storage
22

33
import (
4+
"context"
45
"errors"
6+
"io"
57
"net/http"
8+
"strings"
9+
"sync/atomic"
610
"testing"
711

812
"github.com/minio/minio-go/v7"
13+
"github.com/minio/minio-go/v7/pkg/credentials"
914
"github.com/samber/lo"
1015
"github.com/stretchr/testify/assert"
16+
"github.com/stretchr/testify/require"
1117
)
1218

1319
func TestSplitIntoParts(t *testing.T) {
@@ -89,3 +95,148 @@ func TestIsDeleteSuccessful(t *testing.T) {
8995
})
9096
}
9197
}
98+
99+
func TestBucketExistStopsAfterFirstResult(t *testing.T) {
100+
var listReqCount atomic.Int32
101+
var firstReqPath string
102+
var firstReqMaxKeys string
103+
var firstReqContinuationToken string
104+
cli, err := newInternalMinio(Config{
105+
Provider: "s3",
106+
Endpoint: "example.com",
107+
UseSSL: true,
108+
Bucket: "test-bucket",
109+
Credential: Credential{
110+
Type: Static,
111+
AK: "ak",
112+
SK: "sk",
113+
},
114+
}, &minio.Options{
115+
Secure: true,
116+
Creds: credentials.NewStaticV4("ak", "sk", ""),
117+
Transport: roundTripperFunc(func(r *http.Request) (*http.Response, error) {
118+
if _, ok := r.URL.Query()["location"]; ok {
119+
return &http.Response{
120+
StatusCode: http.StatusOK,
121+
Header: http.Header{"Content-Type": []string{"application/xml"}},
122+
Body: io.NopCloser(strings.NewReader(`<LocationConstraint xmlns="http://s3.amazonaws.com/doc/2006-03-01/"></LocationConstraint>`)),
123+
Request: r,
124+
}, nil
125+
}
126+
127+
if r.URL.Query().Get("list-type") != "2" {
128+
return nil, errors.New("unexpected request")
129+
}
130+
131+
if token := r.URL.Query().Get("continuation-token"); token != "" {
132+
<-r.Context().Done()
133+
return nil, r.Context().Err()
134+
}
135+
136+
listReqCount.Add(1)
137+
firstReqPath = r.URL.Path
138+
firstReqMaxKeys = r.URL.Query().Get("max-keys")
139+
firstReqContinuationToken = r.URL.Query().Get("continuation-token")
140+
141+
return &http.Response{
142+
StatusCode: http.StatusOK,
143+
Header: http.Header{"Content-Type": []string{"application/xml"}},
144+
Body: io.NopCloser(strings.NewReader(`<?xml version="1.0" encoding="UTF-8"?>
145+
<ListBucketResult>
146+
<Name>test-bucket</Name>
147+
<Prefix></Prefix>
148+
<MaxKeys>1</MaxKeys>
149+
<IsTruncated>true</IsTruncated>
150+
<Contents>
151+
<Key>first-object</Key>
152+
<Size>1</Size>
153+
</Contents>
154+
<NextContinuationToken>next-page</NextContinuationToken>
155+
</ListBucketResult>`)),
156+
Request: r,
157+
}, nil
158+
}),
159+
})
160+
require.NoError(t, err)
161+
162+
exists, err := cli.BucketExist(context.Background(), "")
163+
require.NoError(t, err)
164+
assert.True(t, exists)
165+
assert.EqualValues(t, 1, listReqCount.Load())
166+
assert.Equal(t, "/test-bucket/", firstReqPath)
167+
assert.Equal(t, "1", firstReqMaxKeys)
168+
assert.Empty(t, firstReqContinuationToken)
169+
}
170+
171+
func TestBucketExistReturnsFalseForNoSuchBucket(t *testing.T) {
172+
cli, err := newInternalMinio(Config{
173+
Provider: "s3",
174+
Endpoint: "example.com",
175+
UseSSL: true,
176+
Bucket: "missing-bucket",
177+
Credential: Credential{
178+
Type: Static,
179+
AK: "ak",
180+
SK: "sk",
181+
},
182+
}, &minio.Options{
183+
Secure: true,
184+
Creds: credentials.NewStaticV4("ak", "sk", ""),
185+
Transport: roundTripperFunc(func(r *http.Request) (*http.Response, error) {
186+
return &http.Response{
187+
StatusCode: http.StatusNotFound,
188+
Header: http.Header{"Content-Type": []string{"application/xml"}},
189+
Body: io.NopCloser(strings.NewReader(`<?xml version="1.0" encoding="UTF-8"?>
190+
<Error>
191+
<Code>NoSuchBucket</Code>
192+
<Message>The specified bucket does not exist</Message>
193+
<BucketName>missing-bucket</BucketName>
194+
<Resource>/missing-bucket/</Resource>
195+
<RequestId>req</RequestId>
196+
<HostId>host</HostId>
197+
</Error>`)),
198+
Request: r,
199+
}, nil
200+
}),
201+
})
202+
require.NoError(t, err)
203+
204+
exists, err := cli.BucketExist(context.Background(), "")
205+
require.NoError(t, err)
206+
assert.False(t, exists)
207+
}
208+
209+
func TestBucketExistPropagatesContextCancellation(t *testing.T) {
210+
cli, err := newInternalMinio(Config{
211+
Provider: "s3",
212+
Endpoint: "example.com",
213+
UseSSL: true,
214+
Bucket: "test-bucket",
215+
Credential: Credential{
216+
Type: Static,
217+
AK: "ak",
218+
SK: "sk",
219+
},
220+
}, &minio.Options{
221+
Secure: true,
222+
Creds: credentials.NewStaticV4("ak", "sk", ""),
223+
Transport: roundTripperFunc(func(r *http.Request) (*http.Response, error) {
224+
return nil, context.Canceled
225+
}),
226+
})
227+
require.NoError(t, err)
228+
229+
ctx, cancel := context.WithCancel(context.Background())
230+
cancel()
231+
232+
exists, err := cli.BucketExist(ctx, "")
233+
assert.False(t, exists)
234+
require.Error(t, err)
235+
assert.ErrorIs(t, err, context.Canceled)
236+
}
237+
238+
type roundTripperFunc func(*http.Request) (*http.Response, error)
239+
240+
func (f roundTripperFunc) RoundTrip(req *http.Request) (*http.Response, error) {
241+
return f(req)
242+
}

0 commit comments

Comments
 (0)