From 68df7511f63fc9f6fdc45b0ef2b2b13d2b225f2b Mon Sep 17 00:00:00 2001 From: jsas <1351492+jsas@users.noreply.github.com> Date: Fri, 2 Oct 2026 21:34:52 -0700 Subject: [PATCH] filer.remote.sync: do not pin the sync offset on completed work (#11569) * filer.remote.sync: do not pin the sync offset on completed work * filer.remote.sync: a superseded rename uploads the current entry; typed NotFound for a stamp on a deleted entry * filer.remote.sync: a superseded rename keeps the old key when it is the only copy and uploads once * filer.remote.sync: a rename whose content is now remote-only fails the event instead of completing it * filer.remote.sync: a remote-only rename copies the old object to the destination before deleting it * filer.remote.sync: the remote-only rename path follows the filer's current entry and verifies the destination object * filer.remote.sync: an event that described an entry without data is superseded once the filer wrote to it * filer.remote.sync: a superseded rename does only the work left to do uploadCurrentEntry met a remote-only current entry with a fixed error, but a sync plus remote.uncache in the meantime leaves the destination holding the stamped object; that state is complete, not lost. The remote-only case now finishes through completeRemoteOnlyRename, which verifies the destination against the entry stamp and fails only when neither key holds the content. A current entry whose stamp covers its content was already uploaded by the superseding event; skip it instead of writing the same bytes again. * filer.remote.sync: an inherited stamp does not prove the content synced The stamp-coverage skip in uploadCurrentEntry read LastLocalSyncTsNs as proof the current content was uploaded, but a rename carries the source entry's stamp to the destination: a rewrite hidden by that stamp (the case the fallback upload exists for) carries a LastLocalSyncTsNs at or after its mtime and would have been skipped. Drop the check; the remote-only path verifies content at the destination itself through describes. --------- Co-authored-by: James Sas Co-authored-by: Chris Lu --- weed/command/filer_remote_sync_dir.go | 242 +++++++++-- weed/command/filer_remote_sync_dir_test.go | 442 ++++++++++++++++++++- weed/server/filer_grpc_server.go | 3 + weed/util/retry.go | 3 + weed/util/retry_test.go | 4 + 5 files changed, 653 insertions(+), 41 deletions(-) diff --git a/weed/command/filer_remote_sync_dir.go b/weed/command/filer_remote_sync_dir.go index b8051e2cf..7f40a4287 100644 --- a/weed/command/filer_remote_sync_dir.go +++ b/weed/command/filer_remote_sync_dir.go @@ -5,6 +5,7 @@ import ( "context" "errors" "fmt" + "io" "os" "strings" "time" @@ -171,7 +172,6 @@ func (option *RemoteSyncOptions) makeEventProcessor(remoteStorage *remote_pb.Rem glog.V(0).Infof("mkdir %s", remote_storage.FormatLocation(dest)) return client.WriteDirectory(dest, remoteWriteEntry(message.NewEntry, *option.storageClass)) } - glog.V(0).Infof("create %s", remote_storage.FormatLocation(dest)) remoteEntry, writeErr := retriedWriteFile(client, filerSource, message.NewParentPath, remoteWriteEntry(message.NewEntry, *option.storageClass), dest) if errors.Is(writeErr, errSuperseded) { glog.Errorf("skipping %s: %v", remote_storage.FormatLocation(dest), writeErr) @@ -191,20 +191,7 @@ func (option *RemoteSyncOptions) makeEventProcessor(remoteStorage *remote_pb.Rem return updateLocalEntry(option, message.NewParentPath, message.NewEntry, remoteEntry) } if filer_pb.IsDelete(resp) { - // Skip deletion of internal version files; individual version - // deletes should not propagate to the remote object - if isVersionedPath(resp.Directory, message.OldEntry.Name, message.OldEntry.IsDirectory) { - glog.V(2).Infof("skipping delete of internal version path: %s/%s", resp.Directory, message.OldEntry.Name) - return nil - } - glog.V(2).Infof("delete: %+v", resp) - dest := toRemoteStorageLocation(util.FullPath(mountedDir), util.NewFullPath(resp.Directory, message.OldEntry.Name), remoteStorageMountLocation) - if message.OldEntry.IsDirectory { - glog.V(0).Infof("rmdir %s", remote_storage.FormatLocation(dest)) - return client.RemoveDirectory(dest) - } - glog.V(0).Infof("delete %s", remote_storage.FormatLocation(dest)) - return client.DeleteFile(dest) + return processDeleteEvent(client, mountedDir, remoteStorageMountLocation, resp) } if message.OldEntry != nil && message.NewEntry != nil { return processUpdateEvent(option, filerSource, *option.storageClass, client, mountedDir, remoteStorageMountLocation, resp) @@ -215,6 +202,45 @@ func (option *RemoteSyncOptions) makeEventProcessor(remoteStorage *remote_pb.Rem return eachEntryFunc, nil } +// processDeleteEvent removes the remote object an entry mapped to. A remote +// object that is already absent counts as deleted: the entry was never +// uploaded (created and deleted faster than the sync ran, or a replay of an +// event the inline delete already handled), and GCS reports that case as +// ErrRemoteObjectNotFound where S3 and Azure answer an idempotent success. +// Returning the error would pin the sync offset on an event that has nothing +// left to do. +func processDeleteEvent( + client remote_storage.RemoteStorageClient, + mountedDir string, + remoteStorageMountLocation *remote_pb.RemoteStorageLocation, + resp *filer_pb.SubscribeMetadataResponse, +) error { + message := resp.EventNotification + // Skip deletion of internal version files; individual version + // deletes should not propagate to the remote object + if isVersionedPath(resp.Directory, message.OldEntry.Name, message.OldEntry.IsDirectory) { + glog.V(2).Infof("skipping delete of internal version path: %s/%s", resp.Directory, message.OldEntry.Name) + return nil + } + glog.V(2).Infof("delete: %+v", resp) + dest := toRemoteStorageLocation(util.FullPath(mountedDir), util.NewFullPath(resp.Directory, message.OldEntry.Name), remoteStorageMountLocation) + if message.OldEntry.IsDirectory { + glog.V(0).Infof("rmdir %s", remote_storage.FormatLocation(dest)) + return client.RemoveDirectory(dest) + } + glog.V(0).Infof("delete %s", remote_storage.FormatLocation(dest)) + return deleteRemoteFile(client, dest) +} + +// deleteRemoteFile deletes the object and treats an already-absent object as +// deleted. +func deleteRemoteFile(client remote_storage.RemoteStorageClient, dest *remote_pb.RemoteStorageLocation) error { + if err := client.DeleteFile(dest); err != nil && !errors.Is(err, remote_storage.ErrRemoteObjectNotFound) { + return err + } + return nil +} + func processUpdateEvent( filerClient filer_pb.FilerClient, filerSource filer_pb.FilerClient, @@ -256,12 +282,22 @@ func processUpdateEvent( } glog.V(0).Infof("never replicated, uploading %s", remote_storage.FormatLocation(dest)) } - if !proto.Equal(oldDest, dest) && !filer.HasData(message.NewEntry) && message.NewEntry.IsInRemoteOnly() { - glog.V(0).Infof("skip uploading renamed remote-only entry %s: content is only on the deleted remote object", remote_storage.FormatLocation(dest)) - return nil - } glog.V(2).Infof("update: %+v", resp) if !proto.Equal(oldDest, dest) { + // A renamed entry that holds no local data now (the snapshot was + // remote-only, or remote.uncache ran since) has its content only as a + // remote object. Remote-only reads resolve by the entry's own path, so + // the rename completes only once the destination object holds it. The + // filer's current state decides, not the snapshot: an entry rewritten + // since the event has local data again, and the rewrite's bytes are + // what the destination must hold. + current, err := currentEntry(filerSource, message.NewParentPath, message.NewEntry.Name) + if err != nil { + return err + } + if isRemoteOnly(current) { + return completeRemoteOnlyRename(filerClient, client, message.NewParentPath, current, oldDest, dest, storageClass) + } glog.V(0).Infof("delete %s", remote_storage.FormatLocation(oldDest)) if err := client.DeleteFile(oldDest); err != nil { if isMultipartUploadFile(resp.Directory, message.OldEntry.Name) { @@ -275,6 +311,9 @@ func processUpdateEvent( remoteEntry, writeErr := retriedWriteFile(client, filerSource, message.NewParentPath, remoteWriteEntry(message.NewEntry, storageClass), dest) if errors.Is(writeErr, errSuperseded) { glog.Errorf("skipping %s: %v", remote_storage.FormatLocation(dest), writeErr) + if !proto.Equal(oldDest, dest) { + return uploadCurrentEntry(filerClient, filerSource, client, message.NewParentPath, message.NewEntry.Name, oldDest, dest, storageClass) + } return nil } if writeErr != nil { @@ -283,6 +322,128 @@ func processUpdateEvent( return updateLocalEntry(filerClient, message.NewParentPath, message.NewEntry, remoteEntry) } +// currentEntry returns what the filer holds at dir/name now, nil when the +// entry is gone. +func currentEntry(filerSource filer_pb.FilerClient, dir, name string) (*filer_pb.Entry, error) { + current, _, _, err := filer_pb.GetEntry(context.Background(), filerSource, util.NewFullPath(dir, name)) + if errors.Is(err, filer_pb.ErrNotFound) { + return nil, nil + } + if err != nil { + return nil, err + } + return current, nil +} + +// isRemoteOnly reports an entry whose content exists only as its remote +// object: no local data, a RemoteEntry with a size. +func isRemoteOnly(entry *filer_pb.Entry) bool { + return entry != nil && !filer.HasData(entry) && entry.IsInRemoteOnly() +} + +// completeRemoteOnlyRename finishes a rename whose content exists only on the +// remote. When the destination object already holds what the entry's stamp +// describes (the rename ran before, its offset was not persisted, and +// remote.uncache followed), the old key goes. Otherwise the old object is +// copied over the destination, the entry is stamped with the copy, and the +// old key goes. With neither object present the content is lost: the event +// fails and holds the offset for recovery. +func completeRemoteOnlyRename(filerClient filer_pb.FilerClient, client remote_storage.RemoteStorageClient, dir string, current *filer_pb.Entry, oldDest, dest *remote_pb.RemoteStorageLocation, storageClass string) error { + if existing, err := client.StatFile(dest); err == nil { + if describes(current.RemoteEntry, existing) { + glog.V(0).Infof("%s already holds the renamed content", remote_storage.FormatLocation(dest)) + return deleteRemoteFile(client, oldDest) + } + glog.V(0).Infof("%s holds an object the entry does not describe (size %d, want %d); replacing it from %s", remote_storage.FormatLocation(dest), existing.RemoteSize, current.RemoteEntry.RemoteSize, remote_storage.FormatLocation(oldDest)) + } else if !errors.Is(err, remote_storage.ErrRemoteObjectNotFound) { + return err + } + stat, err := client.StatFile(oldDest) + if errors.Is(err, remote_storage.ErrRemoteObjectNotFound) { + return fmt.Errorf("%s: content is on neither %s nor %s", util.NewFullPath(dir, current.Name), remote_storage.FormatLocation(oldDest), remote_storage.FormatLocation(dest)) + } + if err != nil { + return err + } + reader, err := openRemoteObject(client, oldDest, stat.RemoteSize) + if err != nil { + return err + } + defer reader.Close() + glog.V(0).Infof("copy %s -> %s", remote_storage.FormatLocation(oldDest), remote_storage.FormatLocation(dest)) + remoteEntry, err := client.WriteFile(dest, remoteWriteEntry(current, storageClass), reader) + if err != nil { + return err + } + if err := updateLocalEntry(filerClient, dir, current, remoteEntry); err != nil { + return err + } + return deleteRemoteFile(client, oldDest) +} + +// describes reports whether the object a stat returned is the one the entry's +// stamp describes: same size, and the same ETag when both sides carry one (a +// copy written as one stream can legitimately carry a different ETag from a +// multipart original). +func describes(stamp, object *filer_pb.RemoteEntry) bool { + if stamp == nil || object == nil || stamp.RemoteSize != object.RemoteSize { + return false + } + return stamp.RemoteETag == "" || object.RemoteETag == "" || stamp.RemoteETag == object.RemoteETag +} + +// openRemoteObject streams the object when the client can, and reads it whole +// otherwise. +func openRemoteObject(client remote_storage.RemoteStorageClient, loc *remote_pb.RemoteStorageLocation, size int64) (io.ReadCloser, error) { + if streamer, ok := client.(remote_storage.RemoteStorageStreamReader); ok { + return streamer.ReadFileAsStream(context.Background(), loc, 0, size) + } + data, err := client.ReadFile(loc, 0, size) + if err != nil { + return nil, err + } + return io.NopCloser(bytes.NewReader(data)), nil +} + +// uploadCurrentEntry uploads what the filer holds at dir/name now, when the +// event that superseded this rename will not. A rename whose snapshot is +// superseded has already deleted the old key; the rewrite behind it in the +// log uploads the destination itself unless shouldSendToRemote skips it on +// the inherited RemoteEntry, whose RemoteMtime can equal the rewrite's mtime +// within the same second. Only that case uploads here, so the content goes +// up once. A remote-only entry is finished the way completeRemoteOnlyRename +// finishes a rename: its stamp can already describe the destination (a sync +// plus remote.uncache in the meantime), and only then is it complete. +func uploadCurrentEntry(filerClient filer_pb.FilerClient, filerSource filer_pb.FilerClient, client remote_storage.RemoteStorageClient, dir, name string, oldDest, dest *remote_pb.RemoteStorageLocation, storageClass string) error { + current, _, _, err := filer_pb.GetEntry(context.Background(), filerSource, util.NewFullPath(dir, name)) + if errors.Is(err, filer_pb.ErrNotFound) { + return nil + } + if err != nil { + return err + } + if current.IsDirectory { + return nil + } + if isRemoteOnly(current) { + return completeRemoteOnlyRename(filerClient, client, dir, current, oldDest, dest, storageClass) + } + if shouldSendToRemote(current) { + glog.V(0).Infof("leaving %s to the rewrite that superseded the rename", remote_storage.FormatLocation(dest)) + return nil + } + glog.V(0).Infof("uploading the current %s in place of the superseded rename", remote_storage.FormatLocation(dest)) + remoteEntry, writeErr := retriedWriteFile(client, filerSource, dir, remoteWriteEntry(current, storageClass), dest) + if errors.Is(writeErr, errSuperseded) { + glog.Errorf("skipping %s: %v", remote_storage.FormatLocation(dest), writeErr) + return nil + } + if writeErr != nil { + return writeErr + } + return updateLocalEntry(filerClient, dir, current, remoteEntry) +} + // isSuperseded reports whether the filer has moved past the entry an event // described: it is deleted, or it no longer references every chunk the event // named. Those are the chunks the filer deletes when an entry is updated, so @@ -304,6 +465,12 @@ func isSuperseded(filerClient filer_pb.FilerClient, dir string, entry *filer_pb. if err != nil { return false } + if !filer.HasData(entry) && filer.HasData(current) { + // The event described an entry without data (remote-only, or empty); + // the filer has written to it since. Uploading the snapshot would put + // an empty object where the write belongs. + return true + } if len(entry.Content) > 0 || len(current.Content) > 0 { return !bytes.Equal(entry.Content, current.Content) } @@ -314,11 +481,18 @@ func isSuperseded(filerClient filer_pb.FilerClient, dir string, entry *filer_pb. // moved past (isSuperseded). The caller skips the event instead of failing it. var errSuperseded = errors.New("deleted or rewritten since the event was logged") -// retriedWriteFile uploads the entry, retrying transient failures. Every failed -// attempt first asks the filer whether the entry is superseded, and stops at -// once when it is: a dead chunk reads as a transient "RequestError" from the -// SDK, and waiting out the backoff on it buys nothing. +// retriedWriteFile uploads the entry, retrying transient failures. The filer is +// asked whether the entry is superseded before the first attempt and after +// every failed one, and the upload stops at once when it is. Before: a backlog +// (a restart resuming from an old offset) carries every intermediate version +// of a hot file, and uploading each one in turn is wasted bandwidth that can +// trip the remote's per-object mutation rate limit; only the version the filer +// still holds is worth sending. After: a dead chunk reads as a transient +// "RequestError" from the SDK, and waiting out the backoff on it buys nothing. func retriedWriteFile(client remote_storage.RemoteStorageClient, filerSource filer_pb.FilerClient, dir string, newEntry *filer_pb.Entry, dest *remote_pb.RemoteStorageLocation) (remoteEntry *filer_pb.RemoteEntry, err error) { + if isSuperseded(filerSource, dir, newEntry) { + return nil, fmt.Errorf("%s %w", util.NewFullPath(dir, newEntry.Name), errSuperseded) + } err = util.RetryOnError("writeFile", func(err error) bool { return !errors.Is(err, errSuperseded) && util.IsTransientError(err) }, func() error { @@ -453,9 +627,29 @@ func updateLocalEntry(filerClient filer_pb.FilerClient, dir string, entry *filer glog.Errorf("skipping stale stamp of %s: %v", util.NewFullPath(dir, entry.Name), err) return nil } + if isEntryGone(err) { + glog.Errorf("skipping stamp of %s: deleted since the event was logged: %v", util.NewFullPath(dir, entry.Name), err) + return nil + } return err } +// isEntryGone reports an UpdateEntry the filer refused because the entry no +// longer exists: a delete that followed the event superseded the stamp, and +// the delete's own event follows in the log. A filer with the typed answer +// returns codes.NotFound; an older filer returns a plain error of the form +// "not found : ", so the cause (the message's tail, never the +// path) is matched against filer_pb.ErrNotFound's text. +func isEntryGone(err error) bool { + if err == nil { + return false + } + if st, ok := status.FromError(err); ok && st.Code() == codes.NotFound { + return true + } + return strings.HasSuffix(strings.TrimSpace(err.Error()), filer_pb.ErrNotFound.Error()) +} + // ifEntryEqual builds the precondition that the stored entry still equals the // one the event described: chunk fids, inline content, and metadata alike. func ifEntryEqual(entry *filer_pb.Entry) *filer_pb.WriteCondition { @@ -501,7 +695,7 @@ func syncDeleteMarker( dest *remote_pb.RemoteStorageLocation, ) error { glog.V(0).Infof("delete (marker) %s", remote_storage.FormatLocation(dest)) - if err := client.DeleteFile(dest); err != nil { + if err := deleteRemoteFile(client, dest); err != nil { return err } return updateLocalEntry(filerClient, message.NewParentPath, message.NewEntry, &filer_pb.RemoteEntry{ diff --git a/weed/command/filer_remote_sync_dir_test.go b/weed/command/filer_remote_sync_dir_test.go index 7141c2a2c..ade182574 100644 --- a/weed/command/filer_remote_sync_dir_test.go +++ b/weed/command/filer_remote_sync_dir_test.go @@ -415,9 +415,11 @@ func TestIsMetadataOnlyUpdate(t *testing.T) { // every lookup instead. type stubFilerClient struct { filer_pb.SeaweedFilerClient - entry *filer_pb.Entry - err error - lookups int + entry *filer_pb.Entry + err error + updateErr error + lookups int + updates int } func (c *stubFilerClient) LookupDirectoryEntry(context.Context, *filer_pb.LookupDirectoryEntryRequest, ...grpc.CallOption) (*filer_pb.LookupDirectoryEntryResponse, error) { @@ -429,6 +431,10 @@ func (c *stubFilerClient) LookupDirectoryEntry(context.Context, *filer_pb.Lookup } func (c *stubFilerClient) UpdateEntry(context.Context, *filer_pb.UpdateEntryRequest, ...grpc.CallOption) (*filer_pb.UpdateEntryResponse, error) { + c.updates++ + if c.updateErr != nil { + return &filer_pb.UpdateEntryResponse{}, c.updateErr + } return &filer_pb.UpdateEntryResponse{}, nil } @@ -549,6 +555,23 @@ func TestIsSuperseded(t *testing.T) { }) } + t.Run("written since: the event described an entry without data", func(t *testing.T) { + remoteOnly := entryWith("video.mp4", &filer_pb.RemoteEntry{StorageName: "b2", RemoteSize: 20971520, RemoteMtime: 1786096669}) + if !isSuperseded(&stubFilerClient{entry: entryWith("video.mp4", remoteOnly.RemoteEntry, chunk("3,09", "e9"))}, dir, remoteOnly) { + t.Error("isSuperseded = false, want true: the filer wrote local data to a remote-only entry") + } + if isSuperseded(&stubFilerClient{entry: entryWith("video.mp4", remoteOnly.RemoteEntry)}, dir, remoteOnly) { + t.Error("isSuperseded = true, want false: the entry is still remote-only") + } + empty := &filer_pb.Entry{Name: "touch.txt", Attributes: &filer_pb.FuseAttributes{Mtime: 1786096669}} + if !isSuperseded(&stubFilerClient{entry: &filer_pb.Entry{Name: "touch.txt", Content: []byte("now")}}, dir, empty) { + t.Error("isSuperseded = false, want true: the empty file was written to") + } + if isSuperseded(&stubFilerClient{entry: &filer_pb.Entry{Name: "touch.txt"}}, dir, empty) { + t.Error("isSuperseded = true, want false: the file is still empty") + } + }) + t.Run("inline content rewritten since", func(t *testing.T) { event := &filer_pb.Entry{Name: "note.txt", Content: []byte("v1")} if !isSuperseded(&stubFilerClient{entry: &filer_pb.Entry{Name: "note.txt", Content: []byte("v2")}}, dir, event) { @@ -584,7 +607,7 @@ func TestRetriedWriteFileStopsWhenSuperseded(t *testing.T) { // "requesterror" makes IsTransientError retry it to exhaustion. deadChunk := errors.New("RequestError: send request failed\ncaused by: Put \"https://s3.example.com/tier/x\": http://volume:8444/3,01?readDeleted=true: 404 Not Found: not found") - t.Run("superseded, one attempt", func(t *testing.T) { + t.Run("superseded before the first attempt, no write", func(t *testing.T) { remote := &failingRemote{err: deadChunk} filerClient := &stubFilerClient{} start := time.Now() @@ -592,17 +615,32 @@ func TestRetriedWriteFileStopsWhenSuperseded(t *testing.T) { if !errors.Is(err, errSuperseded) { t.Errorf("err = %v, want errSuperseded", err) } + if remote.writes != 0 { + t.Errorf("wrote %d times, want 0: a superseded version is not worth uploading", remote.writes) + } + if filerClient.lookups != 1 { + t.Errorf("looked up the filer %d times, want 1", filerClient.lookups) + } + if elapsed := time.Since(start); elapsed > 500*time.Millisecond { + t.Errorf("took %v, want no backoff", elapsed) + } + }) + + t.Run("superseded during the attempt, one write", func(t *testing.T) { + remote := &failingRemote{err: deadChunk} + filerClient := &supersedingFilerClient{live: entryWith(event.Name, nil, chunk("3,01", "e1"))} + _, err := retriedWriteFile(remote, filerClient, dir, event, dest) + if !errors.Is(err, errSuperseded) { + t.Errorf("err = %v, want errSuperseded", err) + } if !strings.Contains(err.Error(), "404 Not Found") { t.Errorf("err = %v, want it to keep the write failure", err) } if remote.writes != 1 { t.Errorf("wrote %d times, want 1: retrying a dead chunk cannot succeed", remote.writes) } - if filerClient.lookups != 1 { - t.Errorf("looked up the filer %d times, want 1", filerClient.lookups) - } - if elapsed := time.Since(start); elapsed > 500*time.Millisecond { - t.Errorf("took %v, want no backoff", elapsed) + if filerClient.lookups != 2 { + t.Errorf("looked up the filer %d times, want one before the attempt and one after it failed", filerClient.lookups) } }) @@ -619,8 +657,8 @@ func TestRetriedWriteFileStopsWhenSuperseded(t *testing.T) { if remote.writes != 2 { t.Errorf("wrote %d times, want 2: a transient failure on a live entry is still retried", remote.writes) } - if filerClient.lookups != 2 { - t.Errorf("looked up the filer %d times, want one per failed attempt", filerClient.lookups) + if filerClient.lookups != 3 { + t.Errorf("looked up the filer %d times, want one before the first attempt and one per failed attempt", filerClient.lookups) } }) @@ -639,16 +677,51 @@ func TestRetriedWriteFileStopsWhenSuperseded(t *testing.T) { type recordingRemote struct { remote_storage.RemoteStorageClient - deletes []*remote_pb.RemoteStorageLocation - writes []*remote_pb.RemoteStorageLocation - deleteErr error + deletes []*remote_pb.RemoteStorageLocation + writes []*remote_pb.RemoteStorageLocation + written [][]byte + // chunk ids of each written entry: which version of the file went up + writtenChunks [][]string + deleteErr error + // objects present on the remote, by path: StatFile and ReadFile answer + // from it, everything else is ErrRemoteObjectNotFound. + objects map[string][]byte } -func (r *recordingRemote) WriteFile(loc *remote_pb.RemoteStorageLocation, entry *filer_pb.Entry, _ io.Reader) (*filer_pb.RemoteEntry, error) { +func (r *recordingRemote) WriteFile(loc *remote_pb.RemoteStorageLocation, entry *filer_pb.Entry, reader io.Reader) (*filer_pb.RemoteEntry, error) { r.writes = append(r.writes, loc) + var ids []string + for _, c := range entry.GetChunks() { + ids = append(ids, c.GetFileIdString()) + } + r.writtenChunks = append(r.writtenChunks, ids) + // A chunked upload's reader walks volume servers the stub filer cannot + // name, so it is not read; the copy path hands over the old object's + // stream for a chunkless entry, which is drained. + var body []byte + if rc, ok := reader.(io.ReadCloser); ok && len(entry.GetChunks()) == 0 && r.objects != nil { + body, _ = io.ReadAll(rc) + } + r.written = append(r.written, body) return &filer_pb.RemoteEntry{StorageName: loc.Name, RemoteETag: "etag", RemoteSize: int64(len(entry.Content)), RemoteMtime: entry.Attributes.GetMtime()}, nil } +func (r *recordingRemote) StatFile(loc *remote_pb.RemoteStorageLocation) (*filer_pb.RemoteEntry, error) { + body, ok := r.objects[loc.Path] + if !ok { + return nil, remote_storage.ErrRemoteObjectNotFound + } + return &filer_pb.RemoteEntry{StorageName: loc.Name, RemoteSize: int64(len(body))}, nil +} + +func (r *recordingRemote) ReadFile(loc *remote_pb.RemoteStorageLocation, offset int64, size int64) ([]byte, error) { + body, ok := r.objects[loc.Path] + if !ok { + return nil, remote_storage.ErrRemoteObjectNotFound + } + return body[offset : offset+size], nil +} + func (r *recordingRemote) DeleteFile(loc *remote_pb.RemoteStorageLocation) error { r.deletes = append(r.deletes, loc) return r.deleteErr @@ -685,7 +758,8 @@ func TestRenameWithInheritedRemoteEntryWritesNewKey(t *testing.T) { } remote := &recordingRemote{} - filerClient := &stubFilerClient{} + // the filer still holds the renamed entry as the event described it + filerClient := &stubFilerClient{entry: newEntry} if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { t.Fatal(err) } @@ -801,7 +875,8 @@ func TestRenameDeleteOldKeyNotFoundStillWrites(t *testing.T) { } remote := &recordingRemote{deleteErr: remote_storage.ErrRemoteObjectNotFound} - filerClient := &stubFilerClient{} + // the filer still holds the renamed entry as the event described it + filerClient := &stubFilerClient{entry: newEntry} if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { t.Fatalf("err = %v, want nil: an already-deleted old key must not block the write", err) } @@ -880,3 +955,336 @@ func TestIsFailedPrecondition(t *testing.T) { } } } + +// supersedingFilerClient answers the first lookup with the live entry and +// every later one with "not found": the entry is deleted while the upload is +// in flight. +type supersedingFilerClient struct { + filer_pb.SeaweedFilerClient + live *filer_pb.Entry + lookups int +} + +func (c *supersedingFilerClient) LookupDirectoryEntry(context.Context, *filer_pb.LookupDirectoryEntryRequest, ...grpc.CallOption) (*filer_pb.LookupDirectoryEntryResponse, error) { + c.lookups++ + if c.lookups == 1 { + return &filer_pb.LookupDirectoryEntryResponse{Entry: c.live}, nil + } + return nil, filer_pb.ErrNotFound +} + +func (c *supersedingFilerClient) WithFilerClient(_ bool, fn func(filer_pb.SeaweedFilerClient) error) error { + return fn(c) +} + +func (c *supersedingFilerClient) AdjustedUrl(location *filer_pb.Location) string { return location.Url } + +func (c *supersedingFilerClient) GetDataCenter() string { return "" } + +// TestDeleteEventAbsentRemoteObjectIsSuccess: an entry created and deleted +// before the sync uploaded it has no remote object. GCS reports the delete as +// ErrRemoteObjectNotFound; the event has nothing left to do, so it must not +// fail (a failed event pins the sync offset and is replayed on every restart). +func TestDeleteEventAbsentRemoteObjectIsSuccess(t *testing.T) { + const mountedDir = "/buckets" + mountLoc := &remote_pb.RemoteStorageLocation{Name: "gcs", Bucket: "bucket", Path: "/"} + resp := &filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/b/vol", + EventNotification: &filer_pb.EventNotification{ + OldEntry: &filer_pb.Entry{Name: "state.db.tmp", Attributes: &filer_pb.FuseAttributes{Mtime: 1786096669}}, + DeleteChunks: true, + }, + } + wantDelete := &remote_pb.RemoteStorageLocation{Name: "gcs", Bucket: "bucket", Path: "/b/vol/state.db.tmp"} + + t.Run("absent object", func(t *testing.T) { + remote := &recordingRemote{deleteErr: remote_storage.ErrRemoteObjectNotFound} + if err := processDeleteEvent(remote, mountedDir, mountLoc, resp); err != nil { + t.Fatalf("err = %v, want nil: deleting an absent object is complete", err) + } + if len(remote.deletes) != 1 || !proto.Equal(remote.deletes[0], wantDelete) { + t.Errorf("deletes = %+v, want %s", remote.deletes, remote_storage.FormatLocation(wantDelete)) + } + }) + + t.Run("other failure still fails", func(t *testing.T) { + remote := &recordingRemote{deleteErr: errors.New("AccessDenied: Access Denied")} + if err := processDeleteEvent(remote, mountedDir, mountLoc, resp); err == nil { + t.Fatal("err = nil, want the delete failure") + } + }) + + t.Run("delete marker, absent object", func(t *testing.T) { + remote := &recordingRemote{deleteErr: remote_storage.ErrRemoteObjectNotFound} + filerClient := &stubFilerClient{} + markerDelete := &remote_pb.RemoteStorageLocation{Name: "gcs", Bucket: "bucket", Path: "/b/vol/state.db"} + message := &filer_pb.EventNotification{ + NewParentPath: "/buckets/b/vol", + NewEntry: &filer_pb.Entry{Name: "state.db", Attributes: &filer_pb.FuseAttributes{Mtime: 1786096669}}, + } + if err := syncDeleteMarker(remote, filerClient, message, markerDelete); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.deletes) != 1 || !proto.Equal(remote.deletes[0], markerDelete) { + t.Errorf("deletes = %+v, want %s", remote.deletes, remote_storage.FormatLocation(markerDelete)) + } + if filerClient.updates != 1 { + t.Errorf("stamped %d times, want 1: the marker is recorded locally so a replay is a no-op", filerClient.updates) + } + }) +} + +// TestUpdateLocalEntrySkipsDeletedEntry: the stamp that follows an upload lands +// on an entry a later event already deleted. The filer answers "not found"; the +// delete's own event follows in the log, so the stamp is skipped rather than +// failed. +func TestUpdateLocalEntrySkipsDeletedEntry(t *testing.T) { + const dir = "/buckets/b/vol" + remoteEntry := &filer_pb.RemoteEntry{StorageName: "gcs", RemoteSize: 7} + + t.Run("plain not found from the filer", func(t *testing.T) { + filerClient := &stubFilerClient{updateErr: fmt.Errorf("not found %s: %w", dir+"/state.db", filer_pb.ErrNotFound)} + if err := updateLocalEntry(filerClient, dir, entryWith("state.db", nil, chunk("3,01", "e1")), remoteEntry); err != nil { + t.Errorf("err = %v, want nil", err) + } + }) + + t.Run("not found over grpc", func(t *testing.T) { + filerClient := &stubFilerClient{updateErr: status.Error(codes.Unknown, "not found "+dir+"/state.db: "+filer_pb.ErrNotFound.Error())} + if err := updateLocalEntry(filerClient, dir, entryWith("state.db", nil, chunk("3,01", "e1")), remoteEntry); err != nil { + t.Errorf("err = %v, want nil", err) + } + }) + + t.Run("typed not found from the filer", func(t *testing.T) { + filerClient := &stubFilerClient{updateErr: status.Errorf(codes.NotFound, "not found %s/state.db: %v", dir, filer_pb.ErrNotFound)} + if err := updateLocalEntry(filerClient, dir, entryWith("state.db", nil, chunk("3,01", "e1")), remoteEntry); err != nil { + t.Errorf("err = %v, want nil", err) + } + }) + + t.Run("the sentinel inside the path is not a missing entry", func(t *testing.T) { + name := filer_pb.ErrNotFound.Error() + filerClient := &stubFilerClient{updateErr: status.Error(codes.Unknown, "not found "+dir+"/"+name+": database unavailable")} + if err := updateLocalEntry(filerClient, dir, entryWith(name, nil, chunk("3,01", "e1")), remoteEntry); err == nil { + t.Error("err = nil, want the store failure: the stamp must be retried, not dropped") + } + }) + + t.Run("other failure still fails", func(t *testing.T) { + filerClient := &stubFilerClient{updateErr: status.Error(codes.Unavailable, "filer is shutting down")} + if err := updateLocalEntry(filerClient, dir, entryWith("state.db", nil, chunk("3,01", "e1")), remoteEntry); err == nil { + t.Error("err = nil, want the update failure") + } + }) +} + +// TestSupersededRenameUploadsCurrentEntry: a rename A -> B queued behind a +// rewrite of B. The rename's snapshot is superseded (B's chunks changed) and +// the old key A is already deleted. The rewrite's own event uploads B unless +// shouldSendToRemote skips it on the inherited RemoteEntry (same-second +// mtimes); only then does the rename upload what the filer holds now, so the +// content goes up exactly once. +func TestSupersededRenameUploadsCurrentEntry(t *testing.T) { + const mountedDir = "/buckets" + mountLoc := &remote_pb.RemoteStorageLocation{Name: "gcs", Bucket: "bucket", Path: "/"} + eventEntry := entryWith("b.txt", nil, chunk("3,01", "e1")) + resp := &filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/b/dir", + EventNotification: &filer_pb.EventNotification{ + OldEntry: entryWith("a.txt", nil, chunk("3,01", "e1")), + NewParentPath: "/buckets/b/dir", + NewEntry: eventEntry, + }, + } + wantDelete := &remote_pb.RemoteStorageLocation{Name: "gcs", Bucket: "bucket", Path: "/b/dir/a.txt"} + wantWrite := &remote_pb.RemoteStorageLocation{Name: "gcs", Bucket: "bucket", Path: "/b/dir/b.txt"} + + t.Run("rewrite hidden by the inherited stamp: upload once here", func(t *testing.T) { + // RemoteMtime inherited from a.txt equals the rewrite's mtime. + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 1024}, chunk("3,02", "e2")) + remote := &recordingRemote{} + filerClient := &stubFilerClient{entry: current} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.deletes) != 1 || !proto.Equal(remote.deletes[0], wantDelete) { + t.Errorf("deletes = %+v, want the old key %s deleted", remote.deletes, remote_storage.FormatLocation(wantDelete)) + } + if len(remote.writes) != 1 || !proto.Equal(remote.writes[0], wantWrite) { + t.Fatalf("writes = %+v, want the new key %s written once with the current entry", remote.writes, remote_storage.FormatLocation(wantWrite)) + } + if filerClient.updates != 1 { + t.Errorf("stamped %d times, want 1: the current entry is stamped", filerClient.updates) + } + }) + + t.Run("rewrite newer than the stamp: its own event uploads, not this one", func(t *testing.T) { + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096000, RemoteSize: 1024}, chunk("3,02", "e2")) + remote := &recordingRemote{} + filerClient := &stubFilerClient{entry: current} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.deletes) != 1 { + t.Errorf("deletes = %+v, want the old key deleted", remote.deletes) + } + if len(remote.writes) != 0 { + t.Errorf("writes = %+v, want none: the rewrite event behind this one uploads b.txt", remote.writes) + } + }) + + t.Run("deleted meanwhile: nothing to upload", func(t *testing.T) { + remote := &recordingRemote{} + filerClient := &stubFilerClient{} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.writes) != 0 { + t.Errorf("writes = %+v, want none: the delete event follows", remote.writes) + } + }) + + t.Run("uncached meanwhile, old object present: copied to the destination, stamped, old key deleted", func(t *testing.T) { + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 7}) + remote := &recordingRemote{objects: map[string][]byte{"/b/dir/a.txt": []byte("payload")}} + filerClient := &stubFilerClient{entry: current} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.writes) != 1 || !proto.Equal(remote.writes[0], wantWrite) || string(remote.written[0]) != "payload" { + t.Fatalf("writes = %+v (%q), want the old object's bytes written at %s", remote.writes, remote.written, remote_storage.FormatLocation(wantWrite)) + } + if filerClient.updates != 1 { + t.Errorf("stamped %d times, want 1: the entry now points at the destination object", filerClient.updates) + } + if len(remote.deletes) != 1 || !proto.Equal(remote.deletes[0], wantDelete) { + t.Errorf("deletes = %+v, want the old key deleted after the copy", remote.deletes) + } + }) + + t.Run("uncached after an earlier run completed the rename: destination present, old key deleted, nothing written", func(t *testing.T) { + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 7}) + remote := &recordingRemote{objects: map[string][]byte{"/b/dir/a.txt": []byte("payload"), "/b/dir/b.txt": []byte("payload")}} + filerClient := &stubFilerClient{entry: current} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { + t.Fatalf("err = %v, want nil: the replayed rename is already complete", err) + } + if len(remote.writes) != 0 { + t.Errorf("writes = %+v, want none", remote.writes) + } + if len(remote.deletes) != 1 || !proto.Equal(remote.deletes[0], wantDelete) { + t.Errorf("deletes = %+v, want the old key deleted", remote.deletes) + } + }) + + t.Run("uncached and neither object exists: the event fails and holds the offset", func(t *testing.T) { + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 7}) + remote := &recordingRemote{} + filerClient := &stubFilerClient{entry: current} + err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp) + if err == nil || !strings.Contains(err.Error(), "neither") { + t.Fatalf("err = %v, want the lost-content failure", err) + } + if len(remote.writes) != 0 || len(remote.deletes) != 0 { + t.Errorf("writes = %+v deletes = %+v, want none", remote.writes, remote.deletes) + } + }) + + t.Run("destination holds an object of another size: replaced from the old key", func(t *testing.T) { + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 7}) + remote := &recordingRemote{objects: map[string][]byte{"/b/dir/a.txt": []byte("payload"), "/b/dir/b.txt": []byte("something else")}} + filerClient := &stubFilerClient{entry: current} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, resp); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.writes) != 1 || string(remote.written[0]) != "payload" { + t.Fatalf("writes = %+v (%q), want the old object copied over the destination", remote.writes, remote.written) + } + if len(remote.deletes) != 1 || !proto.Equal(remote.deletes[0], wantDelete) { + t.Errorf("deletes = %+v, want the old key deleted after the copy", remote.deletes) + } + }) + + t.Run("snapshot remote-only but rewritten since: the rewrite's bytes go up, nothing is copied", func(t *testing.T) { + remoteOnly := &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 7} + snapshot := &filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/b/dir", + EventNotification: &filer_pb.EventNotification{ + OldEntry: entryWith("a.txt", remoteOnly), + NewParentPath: "/buckets/b/dir", + NewEntry: entryWith("b.txt", remoteOnly), + }, + } + // rewritten with the inherited stamp still hiding it from shouldSendToRemote + current := entryWith("b.txt", remoteOnly, chunk("3,09", "e9")) + remote := &recordingRemote{objects: map[string][]byte{"/b/dir/a.txt": []byte("payload")}} + filerClient := &stubFilerClient{entry: current} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, snapshot); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.writes) != 1 || !proto.Equal(remote.writes[0], wantWrite) { + t.Fatalf("writes = %+v, want one upload of the rewritten b.txt", remote.writes) + } + if len(remote.writtenChunks[0]) != 1 || remote.writtenChunks[0][0] != "3,09" { + t.Errorf("uploaded chunks = %v, want the rewrite's chunk 3,09, not the chunkless snapshot", remote.writtenChunks[0]) + } + if string(remote.written[0]) == "payload" { + t.Error("the old object's bytes were copied over a rewritten entry") + } + if len(remote.deletes) != 1 { + t.Errorf("deletes = %+v, want the old key deleted", remote.deletes) + } + }) + + t.Run("snapshot already remote-only: completed the same way", func(t *testing.T) { + remoteOnly := &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 7} + snapshot := &filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets/b/dir", + EventNotification: &filer_pb.EventNotification{ + OldEntry: entryWith("a.txt", remoteOnly), + NewParentPath: "/buckets/b/dir", + NewEntry: entryWith("b.txt", remoteOnly), + }, + } + remote := &recordingRemote{objects: map[string][]byte{"/b/dir/a.txt": []byte("payload")}} + filerClient := &stubFilerClient{entry: entryWith("b.txt", remoteOnly)} + if err := processUpdateEvent(filerClient, filerClient, "", remote, mountedDir, mountLoc, snapshot); err != nil { + t.Fatalf("err = %v, want nil", err) + } + if len(remote.writes) != 1 || string(remote.written[0]) != "payload" { + t.Errorf("writes = %+v (%q), want the old object copied to the destination", remote.writes, remote.written) + } + if len(remote.deletes) != 1 { + t.Errorf("deletes = %+v, want the old key deleted after the copy", remote.deletes) + } + }) + + t.Run("uncached after the old key was deleted: the event fails and holds the offset", func(t *testing.T) { + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 1024}) + remote := &recordingRemote{} + filerClient := &stubFilerClient{entry: current} + err := uploadCurrentEntry(filerClient, filerClient, remote, "/buckets/b/dir", "b.txt", wantDelete, wantWrite, "") + if err == nil || !strings.Contains(err.Error(), "neither") { + t.Errorf("err = %v, want the lost-content failure", err) + } + if len(remote.writes) != 0 { + t.Errorf("writes = %+v, want none", remote.writes) + } + }) + + t.Run("uncached with the destination already synced: complete without a write", func(t *testing.T) { + current := entryWith("b.txt", &filer_pb.RemoteEntry{StorageName: "gcs", RemoteMtime: 1786096669, RemoteSize: 7}) + remote := &recordingRemote{objects: map[string][]byte{"/b/dir/b.txt": []byte("payload")}} + filerClient := &stubFilerClient{entry: current} + if err := uploadCurrentEntry(filerClient, filerClient, remote, "/buckets/b/dir", "b.txt", wantDelete, wantWrite, ""); err != nil { + t.Fatalf("err = %v, want nil: the destination object matches the entry's stamp", err) + } + if len(remote.writes) != 0 { + t.Errorf("writes = %+v, want none", remote.writes) + } + if len(remote.deletes) != 1 || !proto.Equal(remote.deletes[0], wantDelete) { + t.Errorf("deletes = %+v, want the old key deleted", remote.deletes) + } + }) +} diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index fab4b0a3c..bb1f354ce 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -717,6 +717,9 @@ func (fs *FilerServer) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntr entry, err := fs.filer.FindEntry(ctx, lockPath) if err != nil { + if errors.Is(err, filer_pb.ErrNotFound) { + return &filer_pb.UpdateEntryResponse{}, status.Errorf(codes.NotFound, "not found %s: %v", fullpath, err) + } return &filer_pb.UpdateEntryResponse{}, fmt.Errorf("not found %s: %v", fullpath, err) } if err := validateUpdateEntryPreconditions(entry, req.ExpectedExtended); err != nil { diff --git a/weed/util/retry.go b/weed/util/retry.go index b5017b1d2..03b1fe5d7 100644 --- a/weed/util/retry.go +++ b/weed/util/retry.go @@ -39,6 +39,9 @@ var transientErrorMessages = []string{ "requesttimeout", "slowdown", "throttling", + "ratelimitexceeded", + "rate limit", + "too many requests", "internalerror", "resourceexhausted", "unavailable", diff --git a/weed/util/retry_test.go b/weed/util/retry_test.go index 9bc1d5772..f2782a533 100644 --- a/weed/util/retry_test.go +++ b/weed/util/retry_test.go @@ -83,6 +83,10 @@ func TestIsTransientErrorMessage(t *testing.T) { "Connection reset by peer", "dial tcp 10.0.0.1:8888: connect: no route to host", "rpc error: code = Unavailable desc = the connection is unavailable", + // GCS per-object mutation limit on a hot file; the next attempt after a + // one-second backoff is under it + "googleapi: Error 429: The object bucket/samples.json exceeded the rate limit for object mutation operations (create, update, and delete). Please reduce your request rate. See https://cloud.google.com/storage/docs/gcs429., rateLimitExceeded", + "429 Too Many Requests", } for _, msg := range transient { if !IsTransientErrorMessage(msg) {