diff --git a/weed/operation/upload_content.go b/weed/operation/upload_content.go index 5c02c77ac..4ce3e6149 100644 --- a/weed/operation/upload_content.go +++ b/weed/operation/upload_content.go @@ -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) diff --git a/weed/operation/upload_content_test.go b/weed/operation/upload_content_test.go index 8923d33e5..ebbfb722f 100644 --- a/weed/operation/upload_content_test.go +++ b/weed/operation/upload_content_test.go @@ -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 {