s3: honor assignment fsync in UploadWithRetry (#11449)

* fix: honor assignment fsync in UploadWithRetry

* Tests feedback
This commit is contained in:
Ilia Demianenko
2026-09-26 12:01:52 +08:00
committed by GitHub
parent f7680cf812
commit c58bd0dfd3
2 changed files with 72 additions and 9 deletions
+8 -3
View File
@@ -194,11 +194,12 @@ func NewUploaderWithHttpClient(httpClient HTTPClient) *Uploader {
}
}
func (uploader *Uploader) uploadWithRetryData(assignFn func() (fileId string, host string, auth security.EncodedJwt, err error), uploadOption *UploadOption, data []byte) (fileId string, uploadResult *UploadResult, err error) {
func (uploader *Uploader) uploadWithRetryData(assignFn func() (fileId string, host string, auth security.EncodedJwt, fsync bool, err error), uploadOption *UploadOption, data []byte) (fileId string, uploadResult *UploadResult, err error) {
doUploadFunc := func() error {
var host string
var auth security.EncodedJwt
fileId, host, auth, err = assignFn()
var fsync bool
fileId, host, auth, fsync, err = assignFn()
if err != nil {
return err
}
@@ -208,6 +209,9 @@ func (uploader *Uploader) uploadWithRetryData(assignFn func() (fileId string, ho
genUrl = func(host, fileId string) string { return fmt.Sprintf("http://%s/%s", host, fileId) }
}
uploadOption.UploadUrl = genUrl(host, fileId)
if fsync {
uploadOption.UploadUrl = util_http.AppendQueryParameter(uploadOption.UploadUrl, "fsync", "true")
}
uploadOption.Jwt = auth
if util_http.IsProxyChunkUrl(uploadOption.UploadUrl) {
// The request addresses the filer, which authorizes it and mints the
@@ -263,7 +267,7 @@ func (uploader *Uploader) UploadWithRetry(filerClient filer_pb.FilerClient, assi
uploadOption.Md5 = base64.StdEncoding.EncodeToString(digest[:])
}
fileId, uploadResult, err = uploader.uploadWithRetryData(func() (fileId string, host string, auth security.EncodedJwt, err error) {
fileId, uploadResult, err = uploader.uploadWithRetryData(func() (fileId string, host string, auth security.EncodedJwt, fsync bool, err error) {
// grpc assign volume
if grpcAssignErr := filerClient.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
assignCtx, assignCancel := context.WithTimeout(context.Background(), assignVolumeTimeout)
@@ -278,6 +282,7 @@ func (uploader *Uploader) UploadWithRetry(filerClient filer_pb.FilerClient, assi
}
fileId, auth = resp.FileId, security.EncodedJwt(resp.Auth)
fsync = resp.Fsync
loc := resp.Location
host = filerClient.AdjustedUrl(loc)
+64 -6
View File
@@ -95,12 +95,12 @@ func TestUploadWithRetryDataReassignsOnVolumeSizeExceeded(t *testing.T) {
uploader := newUploader(httpClient)
assignCalls := 0
fileID, uploadResult, err := uploader.uploadWithRetryData(func() (string, string, security.EncodedJwt, error) {
fileID, uploadResult, err := uploader.uploadWithRetryData(func() (string, string, security.EncodedJwt, bool, error) {
assignCalls++
if assignCalls == 1 {
return "1,first", "volume-a", "", nil
return "1,first", "volume-a", "", false, nil
}
return "2,second", "volume-b", "", nil
return "2,second", "volume-b", "", false, nil
}, &UploadOption{Filename: "test.bin"}, []byte("abc"))
if err != nil {
@@ -143,12 +143,12 @@ func TestUploadWithRetryDataReassignsOnReplicaWriteFailure(t *testing.T) {
uploader := newUploader(httpClient)
assignCalls := 0
fileID, uploadResult, err := uploader.uploadWithRetryData(func() (string, string, security.EncodedJwt, error) {
fileID, uploadResult, err := uploader.uploadWithRetryData(func() (string, string, security.EncodedJwt, bool, error) {
assignCalls++
if assignCalls == 1 {
return "1,first", "volume-a", "", nil
return "1,first", "volume-a", "", false, nil
}
return "2,second", "volume-b", "", nil
return "2,second", "volume-b", "", false, nil
}, &UploadOption{Filename: "test.bin"}, []byte("abc"))
if err != nil {
@@ -199,6 +199,64 @@ func (c *bodyCapturingHTTPClient) Do(req *http.Request) (*http.Response, error)
}, nil
}
type fsyncAssignClient struct {
filer_pb.SeaweedFilerClient
response *filer_pb.AssignVolumeResponse
}
func (c *fsyncAssignClient) AssignVolume(context.Context, *filer_pb.AssignVolumeRequest, ...grpc.CallOption) (*filer_pb.AssignVolumeResponse, error) {
return c.response, nil
}
func (c *fsyncAssignClient) WithFilerClient(_ bool, fn func(filer_pb.SeaweedFilerClient) error) error {
return fn(c)
}
func (c *fsyncAssignClient) AdjustedUrl(loc *filer_pb.Location) string { return loc.GetUrl() }
func (c *fsyncAssignClient) GetDataCenter() string { return "" }
func TestUploadWithRetryAppendsFsyncWhenAssigned(t *testing.T) {
for _, tc := range []struct {
name string
proxy bool
}{
{name: "direct"},
{name: "filer proxy", proxy: true},
} {
t.Run(tc.name, func(t *testing.T) {
const fileID = "1,0123456789"
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if gotFsync := r.URL.Query().Get("fsync"); gotFsync != "true" {
t.Errorf("upload fsync = %q, want true when assignment requires fsync", gotFsync)
}
if tc.proxy {
if got := r.URL.Query().Get("proxyChunkId"); got != fileID {
t.Errorf("proxyChunkId = %q, want %q", got, fileID)
}
}
w.WriteHeader(http.StatusCreated)
_, _ = io.WriteString(w, `{"size":3}`)
}))
defer server.Close()
host := strings.TrimPrefix(server.URL, "http://")
client := &fsyncAssignClient{response: &filer_pb.AssignVolumeResponse{
FileId: fileID,
Location: &filer_pb.Location{Url: host},
Fsync: true,
}}
option := &UploadOption{}
if tc.proxy {
option.GenUploadUrl = GenUploadUrlProxy(host)
}
_, _, err, _ := newUploader(server.Client()).UploadWithRetry(client,
&filer_pb.AssignVolumeRequest{Path: "/bucket/object"}, option, strings.NewReader("abc"))
if err != nil {
t.Fatalf("upload failed: %v", err)
}
})
}
}
// hangingAssignSeaweedClient is a SeaweedFilerClient whose AssignVolume blocks
// until its context is done, so a test can prove UploadWithRetry bounds the RPC.
type hangingAssignSeaweedClient struct {