From 96d2d13efedb38908290850a56141a38e3a16c80 Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Wed, 24 Jun 2026 16:31:58 -0700 Subject: [PATCH] s3: replicate by fanning out from the gateway to every holder (#10078) * s3: replicate by fanning out from the gateway to every holder The S3 gateway uploaded each chunk to one volume server, which then relayed the copies to the other replica holders. The gateway now uploads each chunk to every holder in parallel (type=replicate), removing the primary volume server's receive-then-resend relay. AssignVolume returns every replica holder (new repeated Location replicas, forwarded from the master assign), the s3api captures them, and the chunked uploader fans out whenever a chunk has more than one holder. Cipher uploads keep the server-driven path since per-call encryption would diverge the replicas. * s3: cancel sibling replica uploads on the first failure * s3: trim replica fan-out comments * s3: roll back successful fan-out chunk copies when a holder fails A failed fan-out records no FileChunk, so copies that landed on the holders that finished before the cancel were leaked as orphans the caller could not see. Track the holders that succeeded and delete the needle from each (type=replicate, local-only) on failure, leaving nothing behind. --- weed/operation/upload_chunked.go | 133 ++++++++++++--- weed/operation/upload_chunked_test.go | 62 +++++++ weed/pb/filer.proto | 1 + weed/pb/filer_pb/filer.pb.go | 213 +++++++++++++----------- weed/s3api/s3api_object_handlers_put.go | 12 +- weed/server/filer_grpc_server.go | 13 +- 6 files changed, 305 insertions(+), 129 deletions(-) diff --git a/weed/operation/upload_chunked.go b/weed/operation/upload_chunked.go index cfcb07634..84aba9f82 100644 --- a/weed/operation/upload_chunked.go +++ b/weed/operation/upload_chunked.go @@ -15,6 +15,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" "github.com/seaweedfs/seaweedfs/weed/security" + util_http "github.com/seaweedfs/seaweedfs/weed/util/http" ) // ChunkedUploadResult contains the result of a chunked upload @@ -168,9 +169,6 @@ uploadLoop: return } - // Upload chunk data - uploadUrl := fmt.Sprintf("http://%s/%s", assignResult.Url, assignResult.Fid) - // Use per-assignment JWT if present, otherwise fall back to the original JWT // This is critical for secured clusters where each volume assignment has its own JWT jwt := opt.Jwt @@ -182,33 +180,39 @@ uploadLoop: chunkMd5 := md5.Sum(buf.Bytes()) chunkMd5B64 := base64.StdEncoding.EncodeToString(chunkMd5[:]) - uploadOption := &UploadOption{ - UploadUrl: uploadUrl, - Cipher: opt.Cipher, - IsInputCompressed: false, - MimeType: opt.MimeType, - PairMap: nil, - Jwt: jwt, - Md5: chunkMd5B64, - } - var uploadResult *UploadResult var uploadResultErr error - // Use mock upload function if provided (for testing), otherwise use real uploader - if opt.UploadFunc != nil { - uploadResult, uploadResultErr = opt.UploadFunc(ctx, buf.Bytes(), uploadOption) + holders := chunkHolders(assignResult) + // Fan out to every holder, except for cipher: per-call encryption + // would give each replica different bytes, so keep its relay path. + if opt.UploadFunc == nil && !opt.Cipher && len(holders) > 1 { + uploadResult, uploadResultErr = uploadChunkToHolders(ctx, holders, assignResult.Fid, buf.Bytes(), jwt, chunkMd5B64, opt) } else { - uploader, uploaderErr := NewUploader() - if uploaderErr != nil { - uploadErrLock.Lock() - if uploadErr == nil { - uploadErr = fmt.Errorf("create uploader: %w", uploaderErr) + uploadOption := &UploadOption{ + UploadUrl: fmt.Sprintf("http://%s/%s", assignResult.Url, assignResult.Fid), + Cipher: opt.Cipher, + IsInputCompressed: false, + MimeType: opt.MimeType, + PairMap: nil, + Jwt: jwt, + Md5: chunkMd5B64, + } + // Use mock upload function if provided (for testing), otherwise use real uploader + if opt.UploadFunc != nil { + uploadResult, uploadResultErr = opt.UploadFunc(ctx, buf.Bytes(), uploadOption) + } else { + uploader, uploaderErr := NewUploader() + if uploaderErr != nil { + uploadErrLock.Lock() + if uploadErr == nil { + uploadErr = fmt.Errorf("create uploader: %w", uploaderErr) + } + uploadErrLock.Unlock() + return } - uploadErrLock.Unlock() - return + uploadResult, uploadResultErr = uploader.UploadData(ctx, buf.Bytes(), uploadOption) } - uploadResult, uploadResultErr = uploader.UploadData(ctx, buf.Bytes(), uploadOption) } if uploadResultErr != nil { @@ -277,3 +281,84 @@ uploadLoop: SmallContent: nil, }, nil } + +// chunkHolders returns the assigned volume plus its replica holders. +func chunkHolders(assignResult *AssignResult) []string { + hosts := []string{assignResult.Url} + for _, replica := range assignResult.Replicas { + if replica.Url != "" && replica.Url != assignResult.Url { + hosts = append(hosts, replica.Url) + } + } + return hosts +} + +// uploadChunkToHolders writes the chunk to every holder concurrently (each with +// type=replicate). On the first failure it cancels the remaining uploads and +// deletes any copies that already landed, so a partial fan-out leaves no +// orphaned needle the caller cannot see. +func uploadChunkToHolders(ctx context.Context, hosts []string, fid string, data []byte, jwt security.EncodedJwt, md5b64 string, opt *ChunkedUploadOption) (*UploadResult, error) { + uploader, err := NewUploader() + if err != nil { + return nil, fmt.Errorf("create uploader: %w", err) + } + glog.V(4).Infof("replica fan-out: writing chunk %s to %d holders %v", fid, len(hosts), hosts) + ctx, cancel := context.WithCancel(ctx) + defer cancel() + type outcome struct { + host string + result *UploadResult + err error + } + outcomes := make(chan outcome, len(hosts)) + for _, host := range hosts { + go func(host string) { + uploadOption := &UploadOption{ + UploadUrl: fmt.Sprintf("http://%s/%s?type=replicate", host, fid), + Cipher: false, + IsInputCompressed: false, + MimeType: opt.MimeType, + PairMap: nil, + Jwt: jwt, + Md5: md5b64, + } + r, e := uploader.UploadData(ctx, data, uploadOption) + outcomes <- outcome{host, r, e} + }(host) + } + var first *UploadResult + var firstErr error + var succeeded []string + for range hosts { + o := <-outcomes + if o.err != nil { + if firstErr == nil { + firstErr = o.err + cancel() + } + } else { + succeeded = append(succeeded, o.host) + if first == nil { + first = o.result + } + } + } + if firstErr != nil { + // A failed fan-out records no chunk, so roll back the copies that landed + // before the cancel rather than leaking them as orphans. + deleteChunkFromHolders(succeeded, fid, jwt) + return nil, firstErr + } + return first, nil +} + +// deleteChunkFromHolders best-effort removes a needle from each holder it landed +// on, using type=replicate so the volume drops only its local copy. A failed +// delete falls back to vacuum reclaiming the orphan. +func deleteChunkFromHolders(hosts []string, fid string, jwt security.EncodedJwt) { + for _, host := range hosts { + if err := util_http.Delete(fmt.Sprintf("http://%s/%s?type=replicate", host, fid), string(jwt)); err != nil { + glog.Warningf("replica fan-out cleanup: delete %s from %s: %v", fid, host, err) + } + } +} diff --git a/weed/operation/upload_chunked_test.go b/weed/operation/upload_chunked_test.go index 4ea2965bc..0f09baa49 100644 --- a/weed/operation/upload_chunked_test.go +++ b/weed/operation/upload_chunked_test.go @@ -4,8 +4,15 @@ import ( "bytes" "context" "errors" + "fmt" "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + "sync/atomic" "testing" + "time" ) // TestUploadReaderInChunksReturnsPartialResultsOnError verifies that when @@ -229,6 +236,61 @@ func TestUploadReaderInChunksContextCancellation(t *testing.T) { } } +// TestUploadChunkToHoldersRollsBackOnPartialFailure verifies that when a fan-out +// chunk write fails on one holder, the copies that already landed on the other +// holders are deleted (type=replicate, local-only) so nothing is left orphaned. +func TestUploadChunkToHoldersRollsBackOnPartialFailure(t *testing.T) { + const fid = "3,01abcdef" + var goodDeletes, badDeletes int32 + + // Sequence the failure strictly after the good upload so the test does not + // depend on timing: the failing holder returns its error only once the good + // holder has stored the chunk, so the good copy is always what gets rolled back. + goodUploaded := make(chan struct{}) + var once sync.Once + + good := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodDelete { + if strings.Contains(r.URL.Path, "01abcdef") && r.URL.Query().Get("type") == "replicate" { + atomic.AddInt32(&goodDeletes, 1) + } + w.WriteHeader(http.StatusOK) + return + } + fmt.Fprintf(w, `{"name":"f","size":11}`) + once.Do(func() { close(goodUploaded) }) + })) + defer good.Close() + + bad := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodDelete { + atomic.AddInt32(&badDeletes, 1) + w.WriteHeader(http.StatusOK) + return + } + select { + case <-goodUploaded: + case <-time.After(5 * time.Second): + } + w.WriteHeader(http.StatusInternalServerError) + io.WriteString(w, "boom") + })) + defer bad.Close() + + hosts := []string{strings.TrimPrefix(good.URL, "http://"), strings.TrimPrefix(bad.URL, "http://")} + _, err := uploadChunkToHolders(context.Background(), hosts, fid, []byte("hello world"), "", "", &ChunkedUploadOption{}) + + if err == nil { + t.Fatal("expected error from a partial fan-out") + } + if got := atomic.LoadInt32(&goodDeletes); got != 1 { + t.Errorf("expected the succeeded holder to receive 1 cleanup DELETE, got %d", got) + } + if got := atomic.LoadInt32(&badDeletes); got != 0 { + t.Errorf("expected no cleanup DELETE to the failed holder, got %d", got) + } +} + // mockFailingReader simulates a reader that fails after reading some data type mockFailingReader struct { data []byte diff --git a/weed/pb/filer.proto b/weed/pb/filer.proto index 2b068967c..281c36c3f 100644 --- a/weed/pb/filer.proto +++ b/weed/pb/filer.proto @@ -527,6 +527,7 @@ message AssignVolumeResponse { string replication = 7; string error = 8; Location location = 9; + repeated Location replicas = 10; } message LookupVolumeRequest { diff --git a/weed/pb/filer_pb/filer.pb.go b/weed/pb/filer_pb/filer.pb.go index 202bedbe7..928f9d420 100644 --- a/weed/pb/filer_pb/filer.pb.go +++ b/weed/pb/filer_pb/filer.pb.go @@ -1,7 +1,7 @@ // Code generated by protoc-gen-go. DO NOT EDIT. // versions: // protoc-gen-go v1.36.6 -// protoc v6.33.4 +// protoc v7.35.0 // source: filer.proto package filer_pb @@ -3139,6 +3139,7 @@ type AssignVolumeResponse struct { Replication string `protobuf:"bytes,7,opt,name=replication,proto3" json:"replication,omitempty"` Error string `protobuf:"bytes,8,opt,name=error,proto3" json:"error,omitempty"` Location *Location `protobuf:"bytes,9,opt,name=location,proto3" json:"location,omitempty"` + Replicas []*Location `protobuf:"bytes,10,rep,name=replicas,proto3" json:"replicas,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -3222,6 +3223,13 @@ func (x *AssignVolumeResponse) GetLocation() *Location { return nil } +func (x *AssignVolumeResponse) GetReplicas() []*Location { + if x != nil { + return x.Replicas + } + return nil +} + type LookupVolumeRequest struct { state protoimpl.MessageState `protogen:"open.v1"` VolumeIds []string `protobuf:"bytes,1,rep,name=volume_ids,json=volumeIds,proto3" json:"volume_ids,omitempty"` @@ -7081,7 +7089,7 @@ const file_filer_proto_rawDesc = "" + "\tdata_node\x18\t \x01(\tR\bdataNode\x12\x1b\n" + "\tdisk_type\x18\b \x01(\tR\bdiskType\x12,\n" + "\x12expected_data_size\x18\n" + - " \x01(\x04R\x10expectedDataSize\"\xe1\x01\n" + + " \x01(\x04R\x10expectedDataSize\"\x91\x02\n" + "\x14AssignVolumeResponse\x12\x17\n" + "\afile_id\x18\x01 \x01(\tR\x06fileId\x12\x14\n" + "\x05count\x18\x04 \x01(\x05R\x05count\x12\x12\n" + @@ -7091,7 +7099,9 @@ const file_filer_proto_rawDesc = "" + "collection\x12 \n" + "\vreplication\x18\a \x01(\tR\vreplication\x12\x14\n" + "\x05error\x18\b \x01(\tR\x05error\x12.\n" + - "\blocation\x18\t \x01(\v2\x12.filer_pb.LocationR\blocation\"4\n" + + "\blocation\x18\t \x01(\v2\x12.filer_pb.LocationR\blocation\x12.\n" + + "\breplicas\x18\n" + + " \x03(\v2\x12.filer_pb.LocationR\breplicas\"4\n" + "\x13LookupVolumeRequest\x12\x1d\n" + "\n" + "volume_ids\x18\x01 \x03(\tR\tvolumeIds\"=\n" + @@ -7594,104 +7604,105 @@ var file_filer_proto_depIdxs = []int32{ 59, // 36: filer_pb.DeleteEntryResponse.metadata_event:type_name -> filer_pb.SubscribeMetadataResponse 12, // 37: filer_pb.StreamRenameEntryResponse.event_notification:type_name -> filer_pb.EventNotification 45, // 38: filer_pb.AssignVolumeResponse.location:type_name -> filer_pb.Location - 45, // 39: filer_pb.Locations.locations:type_name -> filer_pb.Location - 101, // 40: filer_pb.LookupVolumeResponse.locations_map:type_name -> filer_pb.LookupVolumeResponse.LocationsMapEntry - 47, // 41: filer_pb.CollectionListResponse.collections:type_name -> filer_pb.Collection - 12, // 42: filer_pb.SubscribeMetadataResponse.event_notification:type_name -> filer_pb.EventNotification - 59, // 43: filer_pb.SubscribeMetadataResponse.events:type_name -> filer_pb.SubscribeMetadataResponse - 63, // 44: filer_pb.SubscribeMetadataResponse.log_file_refs:type_name -> filer_pb.LogFileChunkRef - 62, // 45: filer_pb.ListMetadataSubscribersResponse.subscribers:type_name -> filer_pb.MetadataSubscriber - 13, // 46: filer_pb.LogFileChunkRef.chunks:type_name -> filer_pb.FileChunk - 10, // 47: filer_pb.TraverseBfsMetadataResponse.entry:type_name -> filer_pb.Entry - 102, // 48: filer_pb.LocateBrokerResponse.resources:type_name -> filer_pb.LocateBrokerResponse.Resource - 103, // 49: filer_pb.FilerConf.locations:type_name -> filer_pb.FilerConf.PathConf - 10, // 50: filer_pb.CacheRemoteObjectToLocalClusterResponse.entry:type_name -> filer_pb.Entry - 59, // 51: filer_pb.CacheRemoteObjectToLocalClusterResponse.metadata_event:type_name -> filer_pb.SubscribeMetadataResponse - 84, // 52: filer_pb.TransferLocksRequest.locks:type_name -> filer_pb.Lock - 17, // 53: filer_pb.StreamMutateEntryRequest.create_request:type_name -> filer_pb.CreateEntryRequest - 29, // 54: filer_pb.StreamMutateEntryRequest.update_request:type_name -> filer_pb.UpdateEntryRequest - 35, // 55: filer_pb.StreamMutateEntryRequest.delete_request:type_name -> filer_pb.DeleteEntryRequest - 39, // 56: filer_pb.StreamMutateEntryRequest.rename_request:type_name -> filer_pb.StreamRenameEntryRequest - 28, // 57: filer_pb.StreamMutateEntryResponse.create_response:type_name -> filer_pb.CreateEntryResponse - 30, // 58: filer_pb.StreamMutateEntryResponse.update_response:type_name -> filer_pb.UpdateEntryResponse - 36, // 59: filer_pb.StreamMutateEntryResponse.delete_response:type_name -> filer_pb.DeleteEntryResponse - 40, // 60: filer_pb.StreamMutateEntryResponse.rename_response:type_name -> filer_pb.StreamRenameEntryResponse - 95, // 61: filer_pb.MountListResponse.mounts:type_name -> filer_pb.MountInfo - 3, // 62: filer_pb.WriteCondition.Clause.kind:type_name -> filer_pb.WriteCondition.Kind - 44, // 63: filer_pb.LookupVolumeResponse.LocationsMapEntry.value:type_name -> filer_pb.Locations - 5, // 64: filer_pb.SeaweedFiler.LookupDirectoryEntry:input_type -> filer_pb.LookupDirectoryEntryRequest - 7, // 65: filer_pb.SeaweedFiler.ListEntries:input_type -> filer_pb.ListEntriesRequest - 17, // 66: filer_pb.SeaweedFiler.CreateEntry:input_type -> filer_pb.CreateEntryRequest - 29, // 67: filer_pb.SeaweedFiler.UpdateEntry:input_type -> filer_pb.UpdateEntryRequest - 31, // 68: filer_pb.SeaweedFiler.TouchAccessTime:input_type -> filer_pb.TouchAccessTimeRequest - 33, // 69: filer_pb.SeaweedFiler.AppendToEntry:input_type -> filer_pb.AppendToEntryRequest - 35, // 70: filer_pb.SeaweedFiler.DeleteEntry:input_type -> filer_pb.DeleteEntryRequest - 21, // 71: filer_pb.SeaweedFiler.ObjectTransaction:input_type -> filer_pb.ObjectTransactionRequest - 26, // 72: filer_pb.SeaweedFiler.ObjectTransactionBatch:input_type -> filer_pb.ObjectTransactionBatchRequest - 24, // 73: filer_pb.SeaweedFiler.PosixLock:input_type -> filer_pb.PosixLockRequest - 37, // 74: filer_pb.SeaweedFiler.AtomicRenameEntry:input_type -> filer_pb.AtomicRenameEntryRequest - 39, // 75: filer_pb.SeaweedFiler.StreamRenameEntry:input_type -> filer_pb.StreamRenameEntryRequest - 89, // 76: filer_pb.SeaweedFiler.StreamMutateEntry:input_type -> filer_pb.StreamMutateEntryRequest - 41, // 77: filer_pb.SeaweedFiler.AssignVolume:input_type -> filer_pb.AssignVolumeRequest - 43, // 78: filer_pb.SeaweedFiler.LookupVolume:input_type -> filer_pb.LookupVolumeRequest - 48, // 79: filer_pb.SeaweedFiler.CollectionList:input_type -> filer_pb.CollectionListRequest - 50, // 80: filer_pb.SeaweedFiler.DeleteCollection:input_type -> filer_pb.DeleteCollectionRequest - 52, // 81: filer_pb.SeaweedFiler.Statistics:input_type -> filer_pb.StatisticsRequest - 54, // 82: filer_pb.SeaweedFiler.Ping:input_type -> filer_pb.PingRequest - 56, // 83: filer_pb.SeaweedFiler.GetFilerConfiguration:input_type -> filer_pb.GetFilerConfigurationRequest - 64, // 84: filer_pb.SeaweedFiler.TraverseBfsMetadata:input_type -> filer_pb.TraverseBfsMetadataRequest - 58, // 85: filer_pb.SeaweedFiler.SubscribeMetadata:input_type -> filer_pb.SubscribeMetadataRequest - 58, // 86: filer_pb.SeaweedFiler.SubscribeLocalMetadata:input_type -> filer_pb.SubscribeMetadataRequest - 60, // 87: filer_pb.SeaweedFiler.ListMetadataSubscribers:input_type -> filer_pb.ListMetadataSubscribersRequest - 71, // 88: filer_pb.SeaweedFiler.KvGet:input_type -> filer_pb.KvGetRequest - 73, // 89: filer_pb.SeaweedFiler.KvPut:input_type -> filer_pb.KvPutRequest - 76, // 90: filer_pb.SeaweedFiler.CacheRemoteObjectToLocalCluster:input_type -> filer_pb.CacheRemoteObjectToLocalClusterRequest - 78, // 91: filer_pb.SeaweedFiler.DistributedLock:input_type -> filer_pb.LockRequest - 80, // 92: filer_pb.SeaweedFiler.DistributedUnlock:input_type -> filer_pb.UnlockRequest - 82, // 93: filer_pb.SeaweedFiler.FindLockOwner:input_type -> filer_pb.FindLockOwnerRequest - 85, // 94: filer_pb.SeaweedFiler.TransferLocks:input_type -> filer_pb.TransferLocksRequest - 87, // 95: filer_pb.SeaweedFiler.ReplicateLock:input_type -> filer_pb.ReplicateLockRequest - 91, // 96: filer_pb.SeaweedFiler.MountRegister:input_type -> filer_pb.MountRegisterRequest - 93, // 97: filer_pb.SeaweedFiler.MountList:input_type -> filer_pb.MountListRequest - 6, // 98: filer_pb.SeaweedFiler.LookupDirectoryEntry:output_type -> filer_pb.LookupDirectoryEntryResponse - 8, // 99: filer_pb.SeaweedFiler.ListEntries:output_type -> filer_pb.ListEntriesResponse - 28, // 100: filer_pb.SeaweedFiler.CreateEntry:output_type -> filer_pb.CreateEntryResponse - 30, // 101: filer_pb.SeaweedFiler.UpdateEntry:output_type -> filer_pb.UpdateEntryResponse - 32, // 102: filer_pb.SeaweedFiler.TouchAccessTime:output_type -> filer_pb.TouchAccessTimeResponse - 34, // 103: filer_pb.SeaweedFiler.AppendToEntry:output_type -> filer_pb.AppendToEntryResponse - 36, // 104: filer_pb.SeaweedFiler.DeleteEntry:output_type -> filer_pb.DeleteEntryResponse - 22, // 105: filer_pb.SeaweedFiler.ObjectTransaction:output_type -> filer_pb.ObjectTransactionResponse - 27, // 106: filer_pb.SeaweedFiler.ObjectTransactionBatch:output_type -> filer_pb.ObjectTransactionBatchResponse - 25, // 107: filer_pb.SeaweedFiler.PosixLock:output_type -> filer_pb.PosixLockResponse - 38, // 108: filer_pb.SeaweedFiler.AtomicRenameEntry:output_type -> filer_pb.AtomicRenameEntryResponse - 40, // 109: filer_pb.SeaweedFiler.StreamRenameEntry:output_type -> filer_pb.StreamRenameEntryResponse - 90, // 110: filer_pb.SeaweedFiler.StreamMutateEntry:output_type -> filer_pb.StreamMutateEntryResponse - 42, // 111: filer_pb.SeaweedFiler.AssignVolume:output_type -> filer_pb.AssignVolumeResponse - 46, // 112: filer_pb.SeaweedFiler.LookupVolume:output_type -> filer_pb.LookupVolumeResponse - 49, // 113: filer_pb.SeaweedFiler.CollectionList:output_type -> filer_pb.CollectionListResponse - 51, // 114: filer_pb.SeaweedFiler.DeleteCollection:output_type -> filer_pb.DeleteCollectionResponse - 53, // 115: filer_pb.SeaweedFiler.Statistics:output_type -> filer_pb.StatisticsResponse - 55, // 116: filer_pb.SeaweedFiler.Ping:output_type -> filer_pb.PingResponse - 57, // 117: filer_pb.SeaweedFiler.GetFilerConfiguration:output_type -> filer_pb.GetFilerConfigurationResponse - 65, // 118: filer_pb.SeaweedFiler.TraverseBfsMetadata:output_type -> filer_pb.TraverseBfsMetadataResponse - 59, // 119: filer_pb.SeaweedFiler.SubscribeMetadata:output_type -> filer_pb.SubscribeMetadataResponse - 59, // 120: filer_pb.SeaweedFiler.SubscribeLocalMetadata:output_type -> filer_pb.SubscribeMetadataResponse - 61, // 121: filer_pb.SeaweedFiler.ListMetadataSubscribers:output_type -> filer_pb.ListMetadataSubscribersResponse - 72, // 122: filer_pb.SeaweedFiler.KvGet:output_type -> filer_pb.KvGetResponse - 74, // 123: filer_pb.SeaweedFiler.KvPut:output_type -> filer_pb.KvPutResponse - 77, // 124: filer_pb.SeaweedFiler.CacheRemoteObjectToLocalCluster:output_type -> filer_pb.CacheRemoteObjectToLocalClusterResponse - 79, // 125: filer_pb.SeaweedFiler.DistributedLock:output_type -> filer_pb.LockResponse - 81, // 126: filer_pb.SeaweedFiler.DistributedUnlock:output_type -> filer_pb.UnlockResponse - 83, // 127: filer_pb.SeaweedFiler.FindLockOwner:output_type -> filer_pb.FindLockOwnerResponse - 86, // 128: filer_pb.SeaweedFiler.TransferLocks:output_type -> filer_pb.TransferLocksResponse - 88, // 129: filer_pb.SeaweedFiler.ReplicateLock:output_type -> filer_pb.ReplicateLockResponse - 92, // 130: filer_pb.SeaweedFiler.MountRegister:output_type -> filer_pb.MountRegisterResponse - 94, // 131: filer_pb.SeaweedFiler.MountList:output_type -> filer_pb.MountListResponse - 98, // [98:132] is the sub-list for method output_type - 64, // [64:98] is the sub-list for method input_type - 64, // [64:64] is the sub-list for extension type_name - 64, // [64:64] is the sub-list for extension extendee - 0, // [0:64] is the sub-list for field type_name + 45, // 39: filer_pb.AssignVolumeResponse.replicas:type_name -> filer_pb.Location + 45, // 40: filer_pb.Locations.locations:type_name -> filer_pb.Location + 101, // 41: filer_pb.LookupVolumeResponse.locations_map:type_name -> filer_pb.LookupVolumeResponse.LocationsMapEntry + 47, // 42: filer_pb.CollectionListResponse.collections:type_name -> filer_pb.Collection + 12, // 43: filer_pb.SubscribeMetadataResponse.event_notification:type_name -> filer_pb.EventNotification + 59, // 44: filer_pb.SubscribeMetadataResponse.events:type_name -> filer_pb.SubscribeMetadataResponse + 63, // 45: filer_pb.SubscribeMetadataResponse.log_file_refs:type_name -> filer_pb.LogFileChunkRef + 62, // 46: filer_pb.ListMetadataSubscribersResponse.subscribers:type_name -> filer_pb.MetadataSubscriber + 13, // 47: filer_pb.LogFileChunkRef.chunks:type_name -> filer_pb.FileChunk + 10, // 48: filer_pb.TraverseBfsMetadataResponse.entry:type_name -> filer_pb.Entry + 102, // 49: filer_pb.LocateBrokerResponse.resources:type_name -> filer_pb.LocateBrokerResponse.Resource + 103, // 50: filer_pb.FilerConf.locations:type_name -> filer_pb.FilerConf.PathConf + 10, // 51: filer_pb.CacheRemoteObjectToLocalClusterResponse.entry:type_name -> filer_pb.Entry + 59, // 52: filer_pb.CacheRemoteObjectToLocalClusterResponse.metadata_event:type_name -> filer_pb.SubscribeMetadataResponse + 84, // 53: filer_pb.TransferLocksRequest.locks:type_name -> filer_pb.Lock + 17, // 54: filer_pb.StreamMutateEntryRequest.create_request:type_name -> filer_pb.CreateEntryRequest + 29, // 55: filer_pb.StreamMutateEntryRequest.update_request:type_name -> filer_pb.UpdateEntryRequest + 35, // 56: filer_pb.StreamMutateEntryRequest.delete_request:type_name -> filer_pb.DeleteEntryRequest + 39, // 57: filer_pb.StreamMutateEntryRequest.rename_request:type_name -> filer_pb.StreamRenameEntryRequest + 28, // 58: filer_pb.StreamMutateEntryResponse.create_response:type_name -> filer_pb.CreateEntryResponse + 30, // 59: filer_pb.StreamMutateEntryResponse.update_response:type_name -> filer_pb.UpdateEntryResponse + 36, // 60: filer_pb.StreamMutateEntryResponse.delete_response:type_name -> filer_pb.DeleteEntryResponse + 40, // 61: filer_pb.StreamMutateEntryResponse.rename_response:type_name -> filer_pb.StreamRenameEntryResponse + 95, // 62: filer_pb.MountListResponse.mounts:type_name -> filer_pb.MountInfo + 3, // 63: filer_pb.WriteCondition.Clause.kind:type_name -> filer_pb.WriteCondition.Kind + 44, // 64: filer_pb.LookupVolumeResponse.LocationsMapEntry.value:type_name -> filer_pb.Locations + 5, // 65: filer_pb.SeaweedFiler.LookupDirectoryEntry:input_type -> filer_pb.LookupDirectoryEntryRequest + 7, // 66: filer_pb.SeaweedFiler.ListEntries:input_type -> filer_pb.ListEntriesRequest + 17, // 67: filer_pb.SeaweedFiler.CreateEntry:input_type -> filer_pb.CreateEntryRequest + 29, // 68: filer_pb.SeaweedFiler.UpdateEntry:input_type -> filer_pb.UpdateEntryRequest + 31, // 69: filer_pb.SeaweedFiler.TouchAccessTime:input_type -> filer_pb.TouchAccessTimeRequest + 33, // 70: filer_pb.SeaweedFiler.AppendToEntry:input_type -> filer_pb.AppendToEntryRequest + 35, // 71: filer_pb.SeaweedFiler.DeleteEntry:input_type -> filer_pb.DeleteEntryRequest + 21, // 72: filer_pb.SeaweedFiler.ObjectTransaction:input_type -> filer_pb.ObjectTransactionRequest + 26, // 73: filer_pb.SeaweedFiler.ObjectTransactionBatch:input_type -> filer_pb.ObjectTransactionBatchRequest + 24, // 74: filer_pb.SeaweedFiler.PosixLock:input_type -> filer_pb.PosixLockRequest + 37, // 75: filer_pb.SeaweedFiler.AtomicRenameEntry:input_type -> filer_pb.AtomicRenameEntryRequest + 39, // 76: filer_pb.SeaweedFiler.StreamRenameEntry:input_type -> filer_pb.StreamRenameEntryRequest + 89, // 77: filer_pb.SeaweedFiler.StreamMutateEntry:input_type -> filer_pb.StreamMutateEntryRequest + 41, // 78: filer_pb.SeaweedFiler.AssignVolume:input_type -> filer_pb.AssignVolumeRequest + 43, // 79: filer_pb.SeaweedFiler.LookupVolume:input_type -> filer_pb.LookupVolumeRequest + 48, // 80: filer_pb.SeaweedFiler.CollectionList:input_type -> filer_pb.CollectionListRequest + 50, // 81: filer_pb.SeaweedFiler.DeleteCollection:input_type -> filer_pb.DeleteCollectionRequest + 52, // 82: filer_pb.SeaweedFiler.Statistics:input_type -> filer_pb.StatisticsRequest + 54, // 83: filer_pb.SeaweedFiler.Ping:input_type -> filer_pb.PingRequest + 56, // 84: filer_pb.SeaweedFiler.GetFilerConfiguration:input_type -> filer_pb.GetFilerConfigurationRequest + 64, // 85: filer_pb.SeaweedFiler.TraverseBfsMetadata:input_type -> filer_pb.TraverseBfsMetadataRequest + 58, // 86: filer_pb.SeaweedFiler.SubscribeMetadata:input_type -> filer_pb.SubscribeMetadataRequest + 58, // 87: filer_pb.SeaweedFiler.SubscribeLocalMetadata:input_type -> filer_pb.SubscribeMetadataRequest + 60, // 88: filer_pb.SeaweedFiler.ListMetadataSubscribers:input_type -> filer_pb.ListMetadataSubscribersRequest + 71, // 89: filer_pb.SeaweedFiler.KvGet:input_type -> filer_pb.KvGetRequest + 73, // 90: filer_pb.SeaweedFiler.KvPut:input_type -> filer_pb.KvPutRequest + 76, // 91: filer_pb.SeaweedFiler.CacheRemoteObjectToLocalCluster:input_type -> filer_pb.CacheRemoteObjectToLocalClusterRequest + 78, // 92: filer_pb.SeaweedFiler.DistributedLock:input_type -> filer_pb.LockRequest + 80, // 93: filer_pb.SeaweedFiler.DistributedUnlock:input_type -> filer_pb.UnlockRequest + 82, // 94: filer_pb.SeaweedFiler.FindLockOwner:input_type -> filer_pb.FindLockOwnerRequest + 85, // 95: filer_pb.SeaweedFiler.TransferLocks:input_type -> filer_pb.TransferLocksRequest + 87, // 96: filer_pb.SeaweedFiler.ReplicateLock:input_type -> filer_pb.ReplicateLockRequest + 91, // 97: filer_pb.SeaweedFiler.MountRegister:input_type -> filer_pb.MountRegisterRequest + 93, // 98: filer_pb.SeaweedFiler.MountList:input_type -> filer_pb.MountListRequest + 6, // 99: filer_pb.SeaweedFiler.LookupDirectoryEntry:output_type -> filer_pb.LookupDirectoryEntryResponse + 8, // 100: filer_pb.SeaweedFiler.ListEntries:output_type -> filer_pb.ListEntriesResponse + 28, // 101: filer_pb.SeaweedFiler.CreateEntry:output_type -> filer_pb.CreateEntryResponse + 30, // 102: filer_pb.SeaweedFiler.UpdateEntry:output_type -> filer_pb.UpdateEntryResponse + 32, // 103: filer_pb.SeaweedFiler.TouchAccessTime:output_type -> filer_pb.TouchAccessTimeResponse + 34, // 104: filer_pb.SeaweedFiler.AppendToEntry:output_type -> filer_pb.AppendToEntryResponse + 36, // 105: filer_pb.SeaweedFiler.DeleteEntry:output_type -> filer_pb.DeleteEntryResponse + 22, // 106: filer_pb.SeaweedFiler.ObjectTransaction:output_type -> filer_pb.ObjectTransactionResponse + 27, // 107: filer_pb.SeaweedFiler.ObjectTransactionBatch:output_type -> filer_pb.ObjectTransactionBatchResponse + 25, // 108: filer_pb.SeaweedFiler.PosixLock:output_type -> filer_pb.PosixLockResponse + 38, // 109: filer_pb.SeaweedFiler.AtomicRenameEntry:output_type -> filer_pb.AtomicRenameEntryResponse + 40, // 110: filer_pb.SeaweedFiler.StreamRenameEntry:output_type -> filer_pb.StreamRenameEntryResponse + 90, // 111: filer_pb.SeaweedFiler.StreamMutateEntry:output_type -> filer_pb.StreamMutateEntryResponse + 42, // 112: filer_pb.SeaweedFiler.AssignVolume:output_type -> filer_pb.AssignVolumeResponse + 46, // 113: filer_pb.SeaweedFiler.LookupVolume:output_type -> filer_pb.LookupVolumeResponse + 49, // 114: filer_pb.SeaweedFiler.CollectionList:output_type -> filer_pb.CollectionListResponse + 51, // 115: filer_pb.SeaweedFiler.DeleteCollection:output_type -> filer_pb.DeleteCollectionResponse + 53, // 116: filer_pb.SeaweedFiler.Statistics:output_type -> filer_pb.StatisticsResponse + 55, // 117: filer_pb.SeaweedFiler.Ping:output_type -> filer_pb.PingResponse + 57, // 118: filer_pb.SeaweedFiler.GetFilerConfiguration:output_type -> filer_pb.GetFilerConfigurationResponse + 65, // 119: filer_pb.SeaweedFiler.TraverseBfsMetadata:output_type -> filer_pb.TraverseBfsMetadataResponse + 59, // 120: filer_pb.SeaweedFiler.SubscribeMetadata:output_type -> filer_pb.SubscribeMetadataResponse + 59, // 121: filer_pb.SeaweedFiler.SubscribeLocalMetadata:output_type -> filer_pb.SubscribeMetadataResponse + 61, // 122: filer_pb.SeaweedFiler.ListMetadataSubscribers:output_type -> filer_pb.ListMetadataSubscribersResponse + 72, // 123: filer_pb.SeaweedFiler.KvGet:output_type -> filer_pb.KvGetResponse + 74, // 124: filer_pb.SeaweedFiler.KvPut:output_type -> filer_pb.KvPutResponse + 77, // 125: filer_pb.SeaweedFiler.CacheRemoteObjectToLocalCluster:output_type -> filer_pb.CacheRemoteObjectToLocalClusterResponse + 79, // 126: filer_pb.SeaweedFiler.DistributedLock:output_type -> filer_pb.LockResponse + 81, // 127: filer_pb.SeaweedFiler.DistributedUnlock:output_type -> filer_pb.UnlockResponse + 83, // 128: filer_pb.SeaweedFiler.FindLockOwner:output_type -> filer_pb.FindLockOwnerResponse + 86, // 129: filer_pb.SeaweedFiler.TransferLocks:output_type -> filer_pb.TransferLocksResponse + 88, // 130: filer_pb.SeaweedFiler.ReplicateLock:output_type -> filer_pb.ReplicateLockResponse + 92, // 131: filer_pb.SeaweedFiler.MountRegister:output_type -> filer_pb.MountRegisterResponse + 94, // 132: filer_pb.SeaweedFiler.MountList:output_type -> filer_pb.MountListResponse + 99, // [99:133] is the sub-list for method output_type + 65, // [65:99] is the sub-list for method input_type + 65, // [65:65] is the sub-list for extension type_name + 65, // [65:65] is the sub-list for extension extendee + 0, // [0:65] is the sub-list for field type_name } func init() { file_filer_proto_init() } diff --git a/weed/s3api/s3api_object_handlers_put.go b/weed/s3api/s3api_object_handlers_put.go index 45b18527e..9e7736b44 100644 --- a/weed/s3api/s3api_object_handlers_put.go +++ b/weed/s3api/s3api_object_handlers_put.go @@ -503,13 +503,21 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader } // Convert filer_pb.AssignVolumeResponse to operation.AssignResult - return nil, &operation.AssignResult{ + result := &operation.AssignResult{ Fid: assignResult.FileId, Url: assignResult.Location.Url, PublicUrl: assignResult.Location.PublicUrl, Count: uint64(count), Auth: security.EncodedJwt(assignResult.Auth), - }, nil + } + for _, replica := range assignResult.Replicas { + result.Replicas = append(result.Replicas, operation.Location{ + Url: replica.Url, + PublicUrl: replica.PublicUrl, + DataCenter: replica.DataCenter, + }) + } + return nil, result, nil } // Upload with auto-chunking diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index 17a5491d3..3a4c11914 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -759,7 +759,7 @@ func (fs *FilerServer) AssignVolume(ctx context.Context, req *filer_pb.AssignVol return &filer_pb.AssignVolumeResponse{Error: fmt.Sprintf("assign volume result: %v", assignResult.Error)}, nil } - return &filer_pb.AssignVolumeResponse{ + resp = &filer_pb.AssignVolumeResponse{ FileId: assignResult.Fid, Count: int32(assignResult.Count), Location: &filer_pb.Location{ @@ -770,7 +770,16 @@ func (fs *FilerServer) AssignVolume(ctx context.Context, req *filer_pb.AssignVol Auth: string(assignResult.Auth), Collection: so.Collection, Replication: so.Replication, - }, nil + } + // Forward the replica holders so a client can write all copies directly. + for _, replica := range assignResult.Replicas { + resp.Replicas = append(resp.Replicas, &filer_pb.Location{ + Url: replica.Url, + PublicUrl: replica.PublicUrl, + DataCenter: replica.DataCenter, + }) + } + return resp, nil } func (fs *FilerServer) resolveAssignStorageOption(ctx context.Context, req *filer_pb.AssignVolumeRequest) (*operation.StorageOption, error) {