fix: return NoSuchKey and drop stale entries when a remote-mounted object is gone from the remote (#11680)

A remote-only entry whose object was deleted from the remote storage
outside the filer answered GET with 500 and stayed in the filer. The
remote's not-found was lost on the way: the backends' ReadFile returned
it as an untyped error, so FetchAndWriteNeedle failed with codes.Unknown
and nothing downstream could tell it from any other failure.

- remote_storage: GCS, S3 (NoSuchKey) and Azure (BlobNotFound) reads
  return ErrRemoteObjectNotFound. GCS reports a missing bucket the same
  way as a missing object, so it confirms the bucket with a listing.
- volume server and filer: the not-found crosses gRPC as codes.NotFound
  carrying the sentinel's text, and the filer's cache RPC returns
  codes.NotFound, which the S3 gateway already maps to NoSuchKey.
- s3api: the origin fallback answers NoSuchKey on a confirmed not-found.
- filer: a confirmed not-found removes the stale entry, so the lazy
  remote-metadata cache converges on the remote. Only remote-only files
  outside .versions and without an active object lock are removed, only
  if unchanged since the fetch (checked on the object's write owner,
  under the lock S3 object writes take), and with a metadata-only
  delete: the filer skips its inline remote delete, and the delete
  events' entries carry a marker that makes filer.remote.sync and
  filer.remote.gateway skip their remote delete, while filer.sync still
  replicates it. The replicated DeleteEntryRequest carries
  keep_remote_object, so the destination's delete events are marked too.
  The store drops the marker from every write, so clients cannot plant
  it.
This commit is contained in:
Peter Dodd authored and GitHub committed 2026-10-10 11:00:04 +08:00
1 parent d3dc03c85a
commit 5f1a742726
29 files changed
+941 -48

No files matched your search

@@ -287,6 +287,10 @@ func (option *RemoteGatewayOptions) makeBucketedEventProcessor(filerSource *sour
return updateLocalEntry(option, message.NewParentPath, message.NewEntry, remoteEntry)
}
if filer_pb.IsDelete(resp) {
if filer.IsMetadataOnlyDelete(message.OldEntry) {
glog.V(2).Infof("skipping remote delete of metadata-only delete: %s/%s", resp.Directory, message.OldEntry.Name)
return nil
}
if resp.Directory == option.bucketsDir {
return handleDeleteBucket(message.OldEntry)
}
@@ -0,0 +1,53 @@
package command
import (
"testing"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/remote_pb"
"github.com/seaweedfs/seaweedfs/weed/remote_storage"
)
type recordingRemoteMaker struct{ client *recordingRemote }
func (m recordingRemoteMaker) Make(*remote_pb.RemoteConf) (remote_storage.RemoteStorageClient, error) {
return m.client, nil
}
func (m recordingRemoteMaker) HasBucket() bool { return true }
func TestGatewayMetadataOnlyDeleteKeepsTheRemoteObject(t *testing.T) {
remote := &recordingRemote{}
remote_storage.RemoteStorageClientMakers["gatewaytest"] = recordingRemoteMaker{client: remote}
t.Cleanup(func() { delete(remote_storage.RemoteStorageClientMakers, "gatewaytest") })
option := &RemoteGatewayOptions{
bucketsDir: "/buckets",
mappings: &remote_pb.RemoteStorageMapping{Mappings: map[string]*remote_pb.RemoteStorageLocation{"/buckets/b": {Name: "gatewaytest", Bucket: "b", Path: "/"}}},
remoteConfs: map[string]*remote_pb.RemoteConf{"gatewaytest": {Name: "gatewaytest", Type: "gatewaytest"}},
}
process, err := option.makeBucketedEventProcessor(nil)
if err != nil {
t.Fatal(err)
}
deleteEvent := func(extended map[string][]byte) *filer_pb.SubscribeMetadataResponse {
return &filer_pb.SubscribeMetadataResponse{
Directory: "/buckets/b/dir",
EventNotification: &filer_pb.EventNotification{
OldEntry: &filer_pb.Entry{Name: "obj.bin", Attributes: &filer_pb.FuseAttributes{}, Extended: extended},
},
}
}
if err := process(deleteEvent(map[string][]byte{filer.ExtKeepRemoteObjectKey: []byte("true")})); err != nil {
t.Fatal(err)
}
if len(remote.deletes) != 0 {
t.Fatalf("deletes = %+v, want none for a metadata-only delete", remote.deletes)
}
if err := process(deleteEvent(nil)); err != nil {
t.Fatal(err)
}
if len(remote.deletes) != 1 {
t.Fatalf("deletes = %+v, want the ordinary delete", remote.deletes)
}
}
+4
View File
@@ -226,6 +226,10 @@ func processDeleteEvent(
glog.V(2).Infof("skipping delete of internal version path: %s/%s", resp.Directory, message.OldEntry.Name)
return nil
}
if filer.IsMetadataOnlyDelete(message.OldEntry) {
glog.V(2).Infof("skipping remote delete of metadata-only delete: %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 {
@@ -1007,6 +1007,18 @@ func TestDeleteEventAbsentRemoteObjectIsSuccess(t *testing.T) {
}
})
t.Run("metadata-only delete keeps the remote object", func(t *testing.T) {
marked := proto.Clone(resp).(*filer_pb.SubscribeMetadataResponse)
marked.EventNotification.OldEntry.Extended = map[string][]byte{filer.ExtKeepRemoteObjectKey: []byte("true")}
remote := &recordingRemote{}
if err := processDeleteEvent(remote, mountedDir, mountLoc, marked); err != nil {
t.Fatalf("err = %v", err)
}
if len(remote.deletes) != 0 {
t.Errorf("deletes = %+v, want none", remote.deletes)
}
})
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 {
+2 -1
View File
@@ -10,6 +10,7 @@ import (
"sync/atomic"
"time"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/operation"
"github.com/seaweedfs/seaweedfs/weed/pb"
@@ -671,7 +672,7 @@ func genProcessFunction(sourcePath string, targetPath string, excludePaths []str
return nil
}
key := buildKey(dataSink, message, targetPath, sourceOldKey, sourcePath)
return dataSink.DeleteEntry(key, message.OldEntry.IsDirectory, message.DeleteChunks, message.Signatures)
return sink.DeleteEntry(dataSink, key, message.OldEntry.IsDirectory, message.DeleteChunks, filer.IsMetadataOnlyDelete(message.OldEntry), message.Signatures)
}
// handle new entries
+46 -5
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"maps"
"strings"
"time"
@@ -52,6 +53,46 @@ type deleteEntryError struct {
func (e *deleteEntryError) Error() string { return e.msg }
func (e *deleteEntryError) Unwrap() error { return e.cause }
// ExtKeepRemoteObjectKey marks the entry in the event of a metadata-only
// delete, so the remote write-back daemons leave the remote object alone too.
// It is set on the event only; the store drops it from every write, so a
// client cannot plant it on an entry.
const ExtKeepRemoteObjectKey = "Seaweed-X-Keep-Remote-Object"
// IsMetadataOnlyDelete reports a delete event's entry marked by
// WithKeepRemoteObject.
func IsMetadataOnlyDelete(oldEntry *filer_pb.Entry) bool {
_, marked := oldEntry.GetExtended()[ExtKeepRemoteObjectKey]
return marked
}
type keepRemoteObjectKey struct{}
// WithKeepRemoteObject marks a delete as metadata-only: neither the filer nor
// filer.remote.sync deletes the entry's object on the mounted remote storage.
// Unlike isFromOtherCluster, the metadata event stays an ordinary local
// delete, so replication still carries it.
func WithKeepRemoteObject(ctx context.Context) context.Context {
return context.WithValue(ctx, keepRemoteObjectKey{}, true)
}
func keepsRemoteObject(ctx context.Context) bool {
return ctx.Value(keepRemoteObjectKey{}) != nil
}
func deleteEventEntry(ctx context.Context, entry *Entry) *Entry {
if !keepsRemoteObject(ctx) {
return entry
}
marked := entry.ShallowClone()
marked.Extended = maps.Clone(entry.Extended)
if marked.Extended == nil {
marked.Extended = map[string][]byte{}
}
marked.Extended[ExtKeepRemoteObjectKey] = []byte("true")
return marked
}
type OnChunksFunc func([]*filer_pb.FileChunk) error
type OnHardLinkIdsFunc func([]HardLinkId) error
@@ -165,7 +206,7 @@ func (f *Filer) doBatchDeleteFolderMetaAndData(ctx context.Context, entry *Entry
subIsDeletingBucket := f.IsBucket(sub)
err = f.doBatchDeleteFolderMetaAndData(ctx, sub, isRecursive, ignoreRecursiveError, shouldDeleteChunks, subIsDeletingBucket, isFromOtherCluster, nil, onHardLinkIdsFn)
} else {
if !isFromOtherCluster {
if !isFromOtherCluster && !keepsRemoteObject(ctx) {
if _, remoteErr := f.maybeDeleteFromRemote(ctx, sub); remoteErr != nil {
glog.Warningf("remote delete child %s: %v", sub.FullPath, remoteErr)
if !ignoreRecursiveError {
@@ -176,7 +217,7 @@ func (f *Filer) doBatchDeleteFolderMetaAndData(ctx context.Context, entry *Entry
if err != nil && !ignoreRecursiveError {
break
}
f.NotifyUpdateEvent(ctx, sub, nil, shouldDeleteChunks, isFromOtherCluster, nil)
f.NotifyUpdateEvent(ctx, deleteEventEntry(ctx, sub), nil, shouldDeleteChunks, isFromOtherCluster, nil)
if len(sub.HardLinkId) != 0 {
// hard link chunk data are deleted separately
err = onHardLinkIdsFn([]HardLinkId{sub.HardLinkId})
@@ -207,7 +248,7 @@ func (f *Filer) doBatchDeleteFolderMetaAndData(ctx context.Context, entry *Entry
}
}
f.NotifyUpdateEvent(ctx, entry, nil, shouldDeleteChunks, isFromOtherCluster, signatures)
f.NotifyUpdateEvent(ctx, deleteEventEntry(ctx, entry), nil, shouldDeleteChunks, isFromOtherCluster, signatures)
f.DeleteChunks(ctx, entry.FullPath, chunksToDelete)
return nil
@@ -217,7 +258,7 @@ func (f *Filer) doDeleteEntryMetaAndData(ctx context.Context, entry *Entry, shou
glog.V(3).InfofCtx(ctx, "deleting entry %v, delete chunks: %v", entry.FullPath, shouldDeleteChunks)
if !isFromOtherCluster {
if !isFromOtherCluster && !keepsRemoteObject(ctx) {
if _, remoteDeletionErr := f.maybeDeleteFromRemote(ctx, entry); remoteDeletionErr != nil {
return remoteDeletionErr
}
@@ -230,7 +271,7 @@ func (f *Filer) doDeleteEntryMetaAndData(ctx context.Context, entry *Entry, shou
}
if !entry.IsDirectory() {
f.NotifyUpdateEvent(ctx, entry, nil, shouldDeleteChunks, isFromOtherCluster, signatures)
f.NotifyUpdateEvent(ctx, deleteEventEntry(ctx, entry), nil, shouldDeleteChunks, isFromOtherCluster, signatures)
}
return nil
+34
View File
@@ -1,9 +1,13 @@
package filer
import (
"context"
"errors"
"fmt"
"os"
"testing"
"github.com/stretchr/testify/require"
)
// The path in a non-empty-folder failure is the client's, so the marker has to
@@ -27,3 +31,33 @@ func TestNonEmptyFolderClassification(t *testing.T) {
t.Errorf("expected no forgery from the entry name: %v", spoofed)
}
}
func TestKeepRemoteObjectMarkerOnlyLivesOnDeleteEvents(t *testing.T) {
store := newStubFilerStore()
f := newTestFiler(t, store, NewFilerRemoteStorage())
ctx := context.Background()
planted := func() map[string][]byte {
return map[string][]byte{ExtKeepRemoteObjectKey: []byte("true"), "Seaweed-Other": []byte("kept")}
}
dir := &Entry{FullPath: "/buckets/b/dir", Attr: Attr{Mode: os.ModeDir | 0755}}
file := &Entry{FullPath: "/buckets/b/dir/obj.bin", Attr: Attr{Mode: 0644}, Extended: planted()}
require.NoError(t, f.CreateEntry(ctx, dir, nil, false, false, nil, true, 255))
require.NoError(t, f.CreateEntry(ctx, file, nil, false, false, nil, true, 255))
updated := file.ShallowClone()
updated.Extended = planted()
require.NoError(t, f.UpdateEntry(ctx, file, updated, false))
stored := store.entries[string(file.FullPath)]
require.NotContains(t, stored.Extended, ExtKeepRemoteObjectKey, "a client must not be able to store the marker")
require.Contains(t, stored.Extended, "Seaweed-Other")
ordinaryCtx, ordinary := WithMetadataEventSink(ctx)
require.NoError(t, f.DeleteEntryMetaAndData(ordinaryCtx, file.FullPath, false, false, false, false, nil, 0))
require.False(t, IsMetadataOnlyDelete(ordinary.Last().EventNotification.OldEntry))
require.NoError(t, f.CreateEntry(ctx, &Entry{FullPath: file.FullPath, Attr: Attr{Mode: 0644}}, nil, false, false, nil, true, 255))
keepCtx, keep := WithMetadataEventSink(WithKeepRemoteObject(ctx))
require.NoError(t, f.DeleteEntryMetaAndData(keepCtx, dir.FullPath, true, false, false, false, nil, 0))
require.Equal(t, "dir", keep.Last().EventNotification.OldEntry.Name)
require.True(t, IsMetadataOnlyDelete(keep.Last().EventNotification.OldEntry))
}
+3
View File
@@ -141,6 +141,7 @@ func (fsw *FilerStoreWrapper) InsertEntry(ctx context.Context, entry *Entry) err
filer_pb.BeforeEntrySerialization(entry.GetChunks())
normalizeEntryMimeForStore(entry)
delete(entry.Extended, ExtKeepRemoteObjectKey)
if len(entry.HardLinkId) > 0 {
glog.V(4).InfofCtx(ctx, "InsertEntry %s has HardLinkId %x counter=%d",
@@ -170,6 +171,7 @@ func (fsw *FilerStoreWrapper) InsertEntryKnownAbsent(ctx context.Context, entry
filer_pb.BeforeEntrySerialization(entry.GetChunks())
normalizeEntryMimeForStore(entry)
delete(entry.Extended, ExtKeepRemoteObjectKey)
if len(entry.HardLinkId) > 0 {
glog.V(4).InfofCtx(ctx, "InsertEntryKnownAbsent %s has HardLinkId %x counter=%d",
@@ -196,6 +198,7 @@ func (fsw *FilerStoreWrapper) UpdateEntry(ctx context.Context, entry *Entry) err
filer_pb.BeforeEntrySerialization(entry.GetChunks())
normalizeEntryMimeForStore(entry)
delete(entry.Extended, ExtKeepRemoteObjectKey)
if len(entry.HardLinkId) > 0 {
glog.V(4).InfofCtx(ctx, "UpdateEntry %s has HardLinkId %x counter=%d",
+2
View File
@@ -354,6 +354,7 @@ message ObjectMutation {
bytes content = 11; // PATCH_EXTENDED: new Entry.content when set_content
bool touch_mtime = 12; // PATCH_EXTENDED: set the entry's Mtime to now (e.g. a metadata-replace copy)
bool remove_empty_parent = 13; // DELETE: also remove the parent directory when the delete leaves it empty (best-effort)
bool keep_remote_object = 14; // DELETE: leave the object on the mounted remote storage
}
// Recompute re-derives a pointer entry (directory/name on the mutation) from the
@@ -518,6 +519,7 @@ message DeleteEntryRequest {
bool is_from_other_cluster = 7;
repeated int32 signatures = 8;
int64 if_not_modified_after = 9;
bool keep_remote_object = 10;
}
message DeleteEntryResponse {
+23 -4
View File
@@ -1549,6 +1549,7 @@ type ObjectMutation struct {
Content []byte `protobuf:"bytes,11,opt,name=content,proto3" json:"content,omitempty"` // PATCH_EXTENDED: new Entry.content when set_content
TouchMtime bool `protobuf:"varint,12,opt,name=touch_mtime,json=touchMtime,proto3" json:"touch_mtime,omitempty"` // PATCH_EXTENDED: set the entry's Mtime to now (e.g. a metadata-replace copy)
RemoveEmptyParent bool `protobuf:"varint,13,opt,name=remove_empty_parent,json=removeEmptyParent,proto3" json:"remove_empty_parent,omitempty"` // DELETE: also remove the parent directory when the delete leaves it empty (best-effort)
KeepRemoteObject bool `protobuf:"varint,14,opt,name=keep_remote_object,json=keepRemoteObject,proto3" json:"keep_remote_object,omitempty"` // DELETE: leave the object on the mounted remote storage
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -1674,6 +1675,13 @@ func (x *ObjectMutation) GetRemoveEmptyParent() bool {
return false
}
func (x *ObjectMutation) GetKeepRemoteObject() bool {
if x != nil {
return x.KeepRemoteObject
}
return false
}
// Recompute re-derives a pointer entry (directory/name on the mutation) from the
// current contents of a scanned directory, atomically under the transaction's
// lock. It is mechanical: the filer picks the child that sorts first or last by
@@ -2746,6 +2754,7 @@ type DeleteEntryRequest struct {
IsFromOtherCluster bool `protobuf:"varint,7,opt,name=is_from_other_cluster,json=isFromOtherCluster,proto3" json:"is_from_other_cluster,omitempty"`
Signatures []int32 `protobuf:"varint,8,rep,packed,name=signatures,proto3" json:"signatures,omitempty"`
IfNotModifiedAfter int64 `protobuf:"varint,9,opt,name=if_not_modified_after,json=ifNotModifiedAfter,proto3" json:"if_not_modified_after,omitempty"`
KeepRemoteObject bool `protobuf:"varint,10,opt,name=keep_remote_object,json=keepRemoteObject,proto3" json:"keep_remote_object,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -2836,6 +2845,13 @@ func (x *DeleteEntryRequest) GetIfNotModifiedAfter() int64 {
return 0
}
func (x *DeleteEntryRequest) GetKeepRemoteObject() bool {
if x != nil {
return x.KeepRemoteObject
}
return false
}
type DeleteEntryResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
Error string `protobuf:"bytes,1,opt,name=error,proto3" json:"error,omitempty"`
@@ -7362,7 +7378,7 @@ const file_filer_proto_rawDesc = "" +
"\x18IF_EXTENDED_TIME_ELAPSED\x10\b\x12\x13\n" +
"\x0fIF_CHUNKS_EQUAL\x10\t\x12\x12\n" +
"\x0eIF_ENTRY_EQUAL\x10\n" +
"\"\xa2\x05\n" +
"\"\xd0\x05\n" +
"\x0eObjectMutation\x121\n" +
"\x04type\x18\x01 \x01(\x0e2\x1d.filer_pb.ObjectMutation.TypeR\x04type\x12\x1c\n" +
"\tdirectory\x18\x02 \x01(\tR\tdirectory\x12\x12\n" +
@@ -7379,7 +7395,8 @@ const file_filer_proto_rawDesc = "" +
"\acontent\x18\v \x01(\fR\acontent\x12\x1f\n" +
"\vtouch_mtime\x18\f \x01(\bR\n" +
"touchMtime\x12.\n" +
"\x13remove_empty_parent\x18\r \x01(\bR\x11removeEmptyParent\x1a>\n" +
"\x13remove_empty_parent\x18\r \x01(\bR\x11removeEmptyParent\x12,\n" +
"\x12keep_remote_object\x18\x0e \x01(\bR\x10keepRemoteObject\x1a>\n" +
"\x10SetExtendedEntry\x12\x10\n" +
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
"\x05value\x18\x02 \x01(\fR\x05value:\x028\x01\"E\n" +
@@ -7480,7 +7497,7 @@ const file_filer_proto_rawDesc = "" +
"\n" +
"entry_name\x18\x02 \x01(\tR\tentryName\x12+\n" +
"\x06chunks\x18\x03 \x03(\v2\x13.filer_pb.FileChunkR\x06chunks\"\x17\n" +
"\x15AppendToEntryResponse\"\xcb\x02\n" +
"\x15AppendToEntryResponse\"\xf9\x02\n" +
"\x12DeleteEntryRequest\x12\x1c\n" +
"\tdirectory\x18\x01 \x01(\tR\tdirectory\x12\x12\n" +
"\x04name\x18\x02 \x01(\tR\x04name\x12$\n" +
@@ -7491,7 +7508,9 @@ const file_filer_proto_rawDesc = "" +
"\n" +
"signatures\x18\b \x03(\x05R\n" +
"signatures\x121\n" +
"\x15if_not_modified_after\x18\t \x01(\x03R\x12ifNotModifiedAfter\"w\n" +
"\x15if_not_modified_after\x18\t \x01(\x03R\x12ifNotModifiedAfter\x12,\n" +
"\x12keep_remote_object\x18\n" +
" \x01(\bR\x10keepRemoteObject\"w\n" +
"\x13DeleteEntryResponse\x12\x14\n" +
"\x05error\x18\x01 \x01(\tR\x05error\x12J\n" +
"\x0emetadata_event\x18\x02 \x01(\v2#.filer_pb.SubscribeMetadataResponseR\rmetadataEvent\"\xba\x01\n" +
+4
View File
@@ -340,6 +340,10 @@ func DoRemoveWithResponse(ctx context.Context, client SeaweedFilerClient, parent
IsFromOtherCluster: isFromOtherCluster,
Signatures: signatures,
}
return DoRemoveRequest(ctx, client, deleteEntryRequest)
}
func DoRemoveRequest(ctx context.Context, client SeaweedFilerClient, deleteEntryRequest *DeleteEntryRequest) (*DeleteEntryResponse, error) {
if resp, err := client.DeleteEntry(ctx, deleteEntryRequest); err != nil {
if strings.Contains(err.Error(), ErrNotFound.Error()) {
return nil, nil
+66
View File
@@ -1289,6 +1289,16 @@ func (m *ObjectMutation) MarshalToSizedBufferVT(dAtA []byte) (int, error) {
i -= len(m.unknownFields)
copy(dAtA[i:], m.unknownFields)
}
if m.KeepRemoteObject {
i--
if m.KeepRemoteObject {
dAtA[i] = 1
} else {
dAtA[i] = 0
}
i--
dAtA[i] = 0x70
}
if m.RemoveEmptyParent {
i--
if m.RemoveEmptyParent {
@@ -2462,6 +2472,16 @@ func (m *DeleteEntryRequest) MarshalToSizedBufferVT(dAtA []byte) (int, error) {
i -= len(m.unknownFields)
copy(dAtA[i:], m.unknownFields)
}
if m.KeepRemoteObject {
i--
if m.KeepRemoteObject {
dAtA[i] = 1
} else {
dAtA[i] = 0
}
i--
dAtA[i] = 0x50
}
if m.IfNotModifiedAfter != 0 {
i = protohelpers.EncodeVarint(dAtA, i, uint64(m.IfNotModifiedAfter))
i--
@@ -7066,6 +7086,9 @@ func (m *ObjectMutation) SizeVT() (n int) {
if m.RemoveEmptyParent {
n += 2
}
if m.KeepRemoteObject {
n += 2
}
n += len(m.unknownFields)
return n
}
@@ -7495,6 +7518,9 @@ func (m *DeleteEntryRequest) SizeVT() (n int) {
if m.IfNotModifiedAfter != 0 {
n += 1 + protohelpers.SizeOfVarint(uint64(m.IfNotModifiedAfter))
}
if m.KeepRemoteObject {
n += 2
}
n += len(m.unknownFields)
return n
}
@@ -13030,6 +13056,26 @@ func (m *ObjectMutation) UnmarshalVT(dAtA []byte) error {
}
}
m.RemoveEmptyParent = bool(v != 0)
case 14:
if wireType != 0 {
return fmt.Errorf("proto: wrong wireType = %d for field KeepRemoteObject", wireType)
}
var v int
for shift := uint(0); ; shift += 7 {
if shift >= 64 {
return protohelpers.ErrIntOverflow
}
if iNdEx >= l {
return io.ErrUnexpectedEOF
}
b := dAtA[iNdEx]
iNdEx++
v |= int(b&0x7F) << shift
if b < 0x80 {
break
}
}
m.KeepRemoteObject = bool(v != 0)
default:
iNdEx = preIndex
skippy, err := protohelpers.Skip(dAtA[iNdEx:])
@@ -15997,6 +16043,26 @@ func (m *DeleteEntryRequest) UnmarshalVT(dAtA []byte) error {
break
}
}
case 10:
if wireType != 0 {
return fmt.Errorf("proto: wrong wireType = %d for field KeepRemoteObject", wireType)
}
var v int
for shift := uint(0); ; shift += 7 {
if shift >= 64 {
return protohelpers.ErrIntOverflow
}
if iNdEx >= l {
return io.ErrUnexpectedEOF
}
b := dAtA[iNdEx]
iNdEx++
v |= int(b&0x7F) << shift
if b < 0x80 {
break
}
}
m.KeepRemoteObject = bool(v != 0)
default:
iNdEx = preIndex
skippy, err := protohelpers.Skip(dAtA[iNdEx:])
@@ -353,6 +353,9 @@ func (az *azureRemoteStorageClient) ReadFileWithConcurrency(loc *remote_pb.Remot
props, propsErr := blobClient.GetProperties(propsCtx, nil)
cancelProps()
if propsErr != nil {
if bloberror.HasCode(propsErr, bloberror.BlobNotFound) {
return nil, remote_storage.ErrRemoteObjectNotFound
}
return nil, fmt.Errorf("get properties %s%s: %w", loc.Bucket, loc.Path, propsErr)
}
if props.ContentLength == nil {
@@ -374,6 +377,9 @@ func (az *azureRemoteStorageClient) ReadFileWithConcurrency(loc *remote_pb.Remot
Concurrency: uint16(concurrency),
})
if err != nil {
if bloberror.HasCode(err, bloberror.BlobNotFound) {
return nil, remote_storage.ErrRemoteObjectNotFound
}
return nil, fmt.Errorf("failed to download file %s%s: %w", loc.Bucket, loc.Path, err)
}
// Pre-sized buffer: a short read stays zero-padded. Reject it rather than
@@ -3,7 +3,10 @@ package azure
import (
"bytes"
"fmt"
"io"
"net/http"
"os"
"strings"
"testing"
"time"
@@ -430,3 +433,53 @@ func TestAzureErrRemoteObjectNotFoundIsAccessible(t *testing.T) {
require.Error(t, remote_storage.ErrRemoteObjectNotFound)
require.Equal(t, "remote object not found", remote_storage.ErrRemoteObjectNotFound.Error())
}
// azureErrorRoundTripper answers every request with an Azure blob error.
type azureErrorRoundTripper struct {
statusCode int
code string
}
func (e *azureErrorRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
if req.Body != nil {
_, _ = io.Copy(io.Discard, req.Body)
_ = req.Body.Close()
}
body := `<?xml version="1.0" encoding="utf-8"?><Error><Code>` + e.code + `</Code><Message>m</Message></Error>`
return &http.Response{
StatusCode: e.statusCode,
Body: io.NopCloser(strings.NewReader(body)),
Header: http.Header{
"Content-Type": []string{"application/xml"},
"X-Ms-Error-Code": []string{e.code},
},
Request: req,
}, nil
}
func TestAzureReadNotFoundClassification(t *testing.T) {
loc := &remote_pb.RemoteStorageLocation{Name: "test", Bucket: "container", Path: "/obj.bin"}
tests := []struct {
name string
statusCode int
code string
size int64
notFound bool
}{
{"missing blob", http.StatusNotFound, "BlobNotFound", 10, true},
{"missing blob, read to end", http.StatusNotFound, "BlobNotFound", 0, true},
{"missing container", http.StatusNotFound, "ContainerNotFound", 10, false},
{"auth failure", http.StatusForbidden, "AuthenticationFailed", 10, false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
conf := &remote_pb.RemoteConf{Name: "test", AzureAccountName: "testaccount", AzureAccountKey: "dGVzdGtleQ=="}
rs, err := MakeWithHTTPClient(conf, &http.Client{Transport: &azureErrorRoundTripper{statusCode: tt.statusCode, code: tt.code}})
require.NoError(t, err)
_, err = rs.ReadFile(loc, 0, tt.size)
require.Error(t, err)
require.Equal(t, tt.notFound, err == remote_storage.ErrRemoteObjectNotFound, "err: %v", err)
})
}
}
+27 -1
View File
@@ -252,6 +252,9 @@ func (gcs *gcsRemoteStorageClient) ReadFile(loc *remote_pb.RemoteStorageLocation
// breaks range reads and returns sizes that disagree with RemoteSize
rangeReader, readErr := gcs.client.Bucket(loc.Bucket).Object(key).ReadCompressed(true).NewRangeReader(context.Background(), offset, size)
if readErr != nil {
if errors.Is(readErr, storage.ErrObjectNotExist) {
return nil, gcs.missingObjectError(context.Background(), loc)
}
return nil, readErr
}
data, err = io.ReadAll(rangeReader)
@@ -268,13 +271,36 @@ func (gcs *gcsRemoteStorageClient) ReadFileAsStream(ctx context.Context, loc *re
reader, err = gcs.client.Bucket(loc.Bucket).Object(key).ReadCompressed(true).NewRangeReader(ctx, offset, size)
if err != nil {
if errors.Is(err, storage.ErrObjectNotExist) {
return nil, remote_storage.ErrRemoteObjectNotFound
return nil, gcs.missingObjectError(ctx, loc)
}
return nil, fmt.Errorf("failed to open stream for %s%s: %w", loc.Bucket, loc.Path, err)
}
return reader, nil
}
// missingObjectError confirms a read's ErrObjectNotExist before reporting the
// object gone: GCS answers a read from a missing bucket the same way, and that
// must not read as every object under the mount being deleted.
func (gcs *gcsRemoteStorageClient) missingObjectError(ctx context.Context, loc *remote_pb.RemoteStorageLocation) error {
ctx, cancel := context.WithTimeout(ctx, defaultGCSOpTimeout)
defer cancel()
key := loc.Path[1:]
query := &storage.Query{Prefix: key}
if err := query.SetAttrSelection([]string{"Name"}); err != nil {
return fmt.Errorf("read gcs %s%s: %w", loc.Bucket, loc.Path, err)
}
attrs, err := gcs.client.Bucket(loc.Bucket).Objects(ctx, query).Next()
switch {
case errors.Is(err, iterator.Done):
return remote_storage.ErrRemoteObjectNotFound
case err != nil:
return fmt.Errorf("read gcs %s%s: object not found and bucket unconfirmed: %w", loc.Bucket, loc.Path, err)
case attrs.Name == key:
return fmt.Errorf("read gcs %s%s: object not found but listed", loc.Bucket, loc.Path)
}
return remote_storage.ErrRemoteObjectNotFound
}
func (gcs *gcsRemoteStorageClient) WriteDirectory(loc *remote_pb.RemoteStorageLocation, entry *filer_pb.Entry) (err error) {
return nil
}
@@ -1,11 +1,19 @@
package gcs
import (
"context"
"io"
"net/http"
"net/http/httptest"
"strings"
"sync/atomic"
"testing"
"cloud.google.com/go/storage"
"github.com/seaweedfs/seaweedfs/weed/pb/remote_pb"
"github.com/seaweedfs/seaweedfs/weed/remote_storage"
"github.com/stretchr/testify/require"
"google.golang.org/api/option"
)
func TestGCSRemoteStorageClientImplementsInterface(t *testing.T) {
@@ -44,3 +52,58 @@ func TestParseInlineCredentials(t *testing.T) {
_, _, err = ParseInlineCredentials(`not json`)
require.Error(t, err)
}
// fakeGCS answers object reads with readStatus and the JSON object listing
// with listStatus/listBody, counting listings.
type fakeGCS struct {
readStatus int
listStatus int
listBody string
listCalls atomic.Int32
}
func (f *fakeGCS) client(t *testing.T) *gcsRemoteStorageClient {
t.Helper()
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if strings.HasPrefix(r.URL.Path, "/storage/v1/b/") && strings.HasSuffix(r.URL.Path, "/o") {
f.listCalls.Add(1)
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(f.listStatus)
_, _ = io.WriteString(w, f.listBody)
return
}
w.WriteHeader(f.readStatus)
}))
t.Cleanup(srv.Close)
client, err := storage.NewClient(context.Background(), option.WithEndpoint(srv.URL+"/storage/v1/"), option.WithoutAuthentication())
require.NoError(t, err)
t.Cleanup(func() { client.Close() })
return &gcsRemoteStorageClient{conf: &remote_pb.RemoteConf{Name: "test"}, client: client}
}
func TestGCSReadNotFoundClassification(t *testing.T) {
loc := &remote_pb.RemoteStorageLocation{Name: "test", Bucket: "bucket", Path: "/dir/obj.bin"}
tests := []struct {
name string
gcs *fakeGCS
notFound bool
wantLists int32
}{
{"missing object", &fakeGCS{readStatus: http.StatusNotFound, listStatus: http.StatusOK, listBody: `{"items":[{"name":"dir/obj.bin.bak"}]}`}, true, 2},
{"missing bucket", &fakeGCS{readStatus: http.StatusNotFound, listStatus: http.StatusNotFound, listBody: `{"error":{"code":404}}`}, false, 2},
{"bucket listing denied", &fakeGCS{readStatus: http.StatusNotFound, listStatus: http.StatusForbidden, listBody: `{"error":{"code":403}}`}, false, 2},
{"object listed after the read missed it", &fakeGCS{readStatus: http.StatusNotFound, listStatus: http.StatusOK, listBody: `{"items":[{"name":"dir/obj.bin"}]}`}, false, 2},
{"read denied", &fakeGCS{readStatus: http.StatusForbidden}, false, 0},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
client := tt.gcs.client(t)
_, err := client.ReadFile(loc, 0, 10)
require.Equal(t, tt.notFound, err == remote_storage.ErrRemoteObjectNotFound, "ReadFile: %v", err)
_, err = client.ReadFileAsStream(context.Background(), loc, 0, 10)
require.Equal(t, tt.notFound, err == remote_storage.ErrRemoteObjectNotFound, "ReadFileAsStream: %v", err)
require.Equal(t, tt.wantLists, tt.gcs.listCalls.Load())
})
}
}
+1 -1
View File
@@ -80,7 +80,7 @@ type Bucket struct {
CreatedAt time.Time
}
// ErrRemoteObjectNotFound is returned by StatFile when the object does not exist in the remote storage backend.
// ErrRemoteObjectNotFound is returned by StatFile and the ReadFile variants when the object does not exist in the remote storage backend.
var ErrRemoteObjectNotFound = errors.New("remote object not found")
type RemoteStorageClient interface {
+12 -1
View File
@@ -2,6 +2,7 @@ package s3
import (
"context"
"errors"
"fmt"
"io"
"net/http"
@@ -414,6 +415,9 @@ func (s *s3RemoteStorageClient) ReadFileWithConcurrency(loc *remote_pb.RemoteSto
Range: aws.String(fmt.Sprintf("bytes=%d-%d", offset, offset+size-1)),
})
if err != nil {
if isNoSuchKey(err) {
return nil, remote_storage.ErrRemoteObjectNotFound
}
return nil, fmt.Errorf("failed to download file %s%s: %v", loc.Bucket, loc.Path, err)
}
// The buffer is pre-sized to size, so a short read leaves the tail
@@ -433,7 +437,7 @@ func (s *s3RemoteStorageClient) ReadFileAsStream(ctx context.Context, loc *remot
Range: aws.String(fmt.Sprintf("bytes=%d-%d", offset, offset+size-1)),
})
if err != nil {
if aerr, ok := err.(awserr.Error); ok && aerr.Code() == s3.ErrCodeNoSuchKey {
if isNoSuchKey(err) {
return nil, remote_storage.ErrRemoteObjectNotFound
}
return nil, fmt.Errorf("failed to open stream for %s%s: %v", loc.Bucket, loc.Path, err)
@@ -441,6 +445,13 @@ func (s *s3RemoteStorageClient) ReadFileAsStream(ctx context.Context, loc *remot
return output.Body, nil
}
// isNoSuchKey matches only a missing key: a missing bucket is also a 404 but
// says nothing about the object.
func isNoSuchKey(err error) bool {
var aerr awserr.Error
return errors.As(err, &aerr) && aerr.Code() == s3.ErrCodeNoSuchKey
}
func (s *s3RemoteStorageClient) WriteDirectory(loc *remote_pb.RemoteStorageLocation, entry *filer_pb.Entry) (err error) {
return nil
}
@@ -2,6 +2,7 @@ package s3
import (
"bytes"
"context"
"fmt"
"io"
"net/http"
@@ -9,6 +10,7 @@ import (
"testing"
"github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/aws/awserr"
"github.com/aws/aws-sdk-go/aws/credentials"
awss3 "github.com/aws/aws-sdk-go/service/s3"
"github.com/aws/aws-sdk-go/service/s3/s3iface"
@@ -224,9 +226,8 @@ func (c *captureRoundTripper) uploadContentType() string {
return c.uploadReq.Header.Get("Content-Type")
}
func newCapturingS3Client(t *testing.T) (*s3RemoteStorageClient, *captureRoundTripper) {
func newS3ClientWithTransport(t *testing.T, rt http.RoundTripper, supportTagging bool) *s3RemoteStorageClient {
t.Helper()
rt := &captureRoundTripper{}
conf := &remote_pb.RemoteConf{
Type: "s3",
Name: "test",
@@ -235,11 +236,17 @@ func newCapturingS3Client(t *testing.T) (*s3RemoteStorageClient, *captureRoundTr
S3ForcePathStyle: true,
S3AccessKey: "test-key",
S3SecretKey: "test-secret",
S3SupportTagging: supportTagging,
}
httpClient := &http.Client{Transport: rt}
rs, err := MakeWithHTTPClient(conf, httpClient)
rs, err := MakeWithHTTPClient(conf, &http.Client{Transport: rt})
require.NoError(t, err)
return rs.(*s3RemoteStorageClient), rt
return rs.(*s3RemoteStorageClient)
}
func newCapturingS3Client(t *testing.T) (*s3RemoteStorageClient, *captureRoundTripper) {
t.Helper()
rt := &captureRoundTripper{}
return newS3ClientWithTransport(t, rt, false), rt
}
func TestS3WriteFilePassesMimeAsContentType(t *testing.T) {
@@ -281,19 +288,7 @@ func (c *recordingRoundTripper) RoundTrip(req *http.Request) (*http.Response, er
func newRecordingS3Client(t *testing.T, supportTagging bool) (*s3RemoteStorageClient, *recordingRoundTripper) {
t.Helper()
rt := &recordingRoundTripper{}
conf := &remote_pb.RemoteConf{
Type: "s3",
Name: "test",
S3Region: "us-east-1",
S3Endpoint: "https://example.invalid",
S3ForcePathStyle: true,
S3AccessKey: "test-key",
S3SecretKey: "test-secret",
S3SupportTagging: supportTagging,
}
rs, err := MakeWithHTTPClient(conf, &http.Client{Transport: rt})
require.NoError(t, err)
return rs.(*s3RemoteStorageClient), rt
return newS3ClientWithTransport(t, rt, supportTagging), rt
}
func TestS3UpdateFileMetadataSkipsTaggingWhenUnsupported(t *testing.T) {
@@ -345,3 +340,55 @@ func TestS3WriteFileOmitsContentTypeWhenMimeMissing(t *testing.T) {
// remote can apply its own default rather than getting a misleading one.
require.Equal(t, "", rt.uploadContentType())
}
// errorRoundTripper answers every request with an S3 XML error.
type errorRoundTripper struct {
statusCode int
code string
}
func (e *errorRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
if req.Body != nil {
_, _ = io.Copy(io.Discard, req.Body)
_ = req.Body.Close()
}
body := `<?xml version="1.0" encoding="UTF-8"?><Error><Code>` + e.code + `</Code><Message>m</Message></Error>`
return &http.Response{
StatusCode: e.statusCode,
Body: io.NopCloser(strings.NewReader(body)),
Header: http.Header{"Content-Type": []string{"application/xml"}},
Request: req,
}, nil
}
func TestS3ReadNotFoundClassification(t *testing.T) {
loc := &remote_pb.RemoteStorageLocation{Name: "test", Bucket: "bucket", Path: "/obj.bin"}
tests := []struct {
name string
statusCode int
code string
notFound bool
}{
{"missing key", http.StatusNotFound, "NoSuchKey", true},
{"missing bucket", http.StatusNotFound, "NoSuchBucket", false},
{"access denied", http.StatusForbidden, "AccessDenied", false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
client := newS3ClientWithTransport(t, &errorRoundTripper{statusCode: tt.statusCode, code: tt.code}, false)
_, err := client.ReadFile(loc, 0, 10)
require.Error(t, err)
require.Equal(t, tt.notFound, err == remote_storage.ErrRemoteObjectNotFound, "ReadFile: %v", err)
_, err = client.ReadFileAsStream(context.Background(), loc, 0, 10)
require.Error(t, err)
require.Equal(t, tt.notFound, err == remote_storage.ErrRemoteObjectNotFound, "ReadFileAsStream: %v", err)
})
}
}
func TestIsNoSuchKey(t *testing.T) {
require.True(t, isNoSuchKey(fmt.Errorf("download: %w", awserr.New(awss3.ErrCodeNoSuchKey, "missing", nil))))
require.False(t, isNoSuchKey(awserr.New(awss3.ErrCodeNoSuchBucket, "missing", nil)))
}
+2 -1
View File
@@ -5,6 +5,7 @@ import (
"fmt"
"time"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
@@ -100,7 +101,7 @@ func (r *Replicator) Replicate(ctx context.Context, key string, message *filer_p
if oldEntry != nil && newEntry == nil {
glog.V(4).Infof("deleting %v", oldSinkKey)
return r.sink.DeleteEntry(oldSinkKey, oldEntry.IsDirectory, message.DeleteChunks, message.Signatures)
return sink.DeleteEntry(r.sink, oldSinkKey, oldEntry.IsDirectory, message.DeleteChunks, filer.IsMetadataOnlyDelete(oldEntry), message.Signatures)
}
if oldEntry == nil && newEntry != nil {
glog.V(4).Infof("creating %v", oldSinkKey)
+35
View File
@@ -4,6 +4,7 @@ import (
"context"
"testing"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/replication/sink"
"github.com/seaweedfs/seaweedfs/weed/replication/source"
@@ -409,3 +410,37 @@ func TestReplicateRenameFromExcludedDirBecomesCreate(t *testing.T) {
t.Fatalf("unexpected delete/update calls: deletes=%+v updates=%+v", s.deleteCalls, s.updateCalls)
}
}
type keepingSink struct {
*recordingSink
keepCalls []deleteCall
}
var _ sink.MetadataOnlyDeleter = (*keepingSink)(nil)
func (s *keepingSink) DeleteEntryKeepingRemoteObject(key string, isDirectory, deleteIncludeChunks bool, signatures []int32) error {
s.keepCalls = append(s.keepCalls, deleteCall{key: key, isDirectory: isDirectory})
return nil
}
func TestReplicateMetadataOnlyDeleteKeepsTheRemoteObject(t *testing.T) {
s := &keepingSink{recordingSink: &recordingSink{name: "filer", sinkToDirectory: "/dest"}}
r := &Replicator{sink: s, source: &source.FilerSource{Dir: "/source"}}
deleteOf := func(name string, extended map[string][]byte) *filer_pb.EventNotification {
return &filer_pb.EventNotification{OldEntry: &filer_pb.Entry{Name: name, Attributes: &filer_pb.FuseAttributes{}, Extended: extended}}
}
if err := r.Replicate(context.Background(), "/source/dir/kept.bin", deleteOf("kept.bin", map[string][]byte{filer.ExtKeepRemoteObjectKey: []byte("true")})); err != nil {
t.Fatal(err)
}
if err := r.Replicate(context.Background(), "/source/dir/gone.bin", deleteOf("gone.bin", nil)); err != nil {
t.Fatal(err)
}
if len(s.keepCalls) != 1 || s.keepCalls[0].key != "/dest/dir/kept.bin" {
t.Fatalf("metadata-only deletes = %+v, want /dest/dir/kept.bin", s.keepCalls)
}
if len(s.deleteCalls) != 1 || s.deleteCalls[0].key != "/dest/dir/gone.bin" {
t.Fatalf("ordinary deletes = %+v, want /dest/dir/gone.bin", s.deleteCalls)
}
}
+21 -1
View File
@@ -172,11 +172,31 @@ func (fs *FilerSink) ActiveTransfers() []ChunkTransferSnapshot {
}
func (fs *FilerSink) DeleteEntry(key string, isDirectory, deleteIncludeChunks bool, signatures []int32) error {
return fs.deleteEntry(key, deleteIncludeChunks, false, signatures)
}
func (fs *FilerSink) DeleteEntryKeepingRemoteObject(key string, isDirectory, deleteIncludeChunks bool, signatures []int32) error {
return fs.deleteEntry(key, deleteIncludeChunks, true, signatures)
}
func (fs *FilerSink) deleteEntry(key string, deleteIncludeChunks, keepRemoteObject bool, signatures []int32) error {
dir, name := util.FullPath(key).DirAndName()
glog.V(4).Infof("delete entry: %v", key)
err := filer_pb.Remove(context.Background(), fs, dir, name, deleteIncludeChunks, true, true, true, signatures)
err := fs.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
_, err := filer_pb.DoRemoveRequest(context.Background(), client, &filer_pb.DeleteEntryRequest{
Directory: dir,
Name: name,
IsDeleteData: deleteIncludeChunks,
IsRecursive: true,
IgnoreRecursiveError: true,
IsFromOtherCluster: true,
Signatures: signatures,
KeepRemoteObject: keepRemoteObject,
})
return err
})
if err != nil {
glog.V(0).Infof("delete entry %s: %v", key, err)
return fmt.Errorf("delete entry %s: %w", key, err)
+15
View File
@@ -32,6 +32,21 @@ type EntryMover interface {
MoveEntry(oldKey, newKey string, newEntry *filer_pb.Entry, signatures []int32) error
}
// MetadataOnlyDeleter is an optional capability for sinks whose own delete
// events reach remote write-back daemons: the destination filer re-emits a
// replicated delete, so a source delete that kept the remote object must stay
// metadata-only there too.
type MetadataOnlyDeleter interface {
DeleteEntryKeepingRemoteObject(key string, isDirectory, deleteIncludeChunks bool, signatures []int32) error
}
func DeleteEntry(s ReplicationSink, key string, isDirectory, deleteIncludeChunks, keepRemoteObject bool, signatures []int32) error {
if deleter, ok := s.(MetadataOnlyDeleter); ok && keepRemoteObject {
return deleter.DeleteEntryKeepingRemoteObject(key, isDirectory, deleteIncludeChunks, signatures)
}
return s.DeleteEntry(key, isDirectory, deleteIncludeChunks, signatures)
}
var (
Sinks []ReplicationSink
)
+16 -7
View File
@@ -1169,10 +1169,8 @@ func (s3a *S3ApiServer) streamFromVolumeServers(w http.ResponseWriter, r *http.R
entry = cachedEntry
glog.V(1).Infof("streamFromVolumeServers: successfully cached remote object, got %d chunks", len(chunks))
} else if isFilerNotFound(cacheErr) {
// Authoritative: the entry vanished; the origin cannot resurrect it
glog.Errorf("streamFromVolumeServers: entry not found while caching %s/%s: %v", bucket, object, cacheErr)
s3err.WriteErrorResponse(w, r, s3err.ErrNoSuchKey)
return newStreamErrorWithResponse(cacheErr)
// Authoritative: the entry vanished or its remote object is gone
return respondObjectGone(w, r, bucket, object, cacheErr)
} else {
// Client disconnected during the cache wait: report cancellation, not an
// error response, so we don't write to a closed connection.
@@ -1185,12 +1183,16 @@ func (s3a *S3ApiServer) streamFromVolumeServers(w http.ResponseWriter, r *http.R
// latest-version read -- has no origin key, so it keeps the error
// paths below.
if cacheVersionId == "" || cacheVersionId == "null" {
if served, streamErr := s3a.serveObjectFromRemoteMount(w, r, entry, bucket, object, offset, size, isRangeRequest, totalSize, t0); served {
served, streamErr := s3a.serveObjectFromRemoteMount(w, r, entry, bucket, object, offset, size, isRangeRequest, totalSize, t0)
if served {
if streamErr != nil {
return newStreamErrorWithResponse(streamErr)
}
return nil
}
if errors.Is(streamErr, remote_storage.ErrRemoteObjectNotFound) {
return respondObjectGone(w, r, bucket, object, streamErr)
}
}
// Origin unreadable. A permanent cache error is final; a transient one
// gets 503 so the client retries the still-filling cache.
@@ -1398,14 +1400,21 @@ func probeReadable(ctx context.Context, reader ctxReaderAt, offset int64, timeou
return err
}
func respondObjectGone(w http.ResponseWriter, r *http.Request, bucket, object string, err error) error {
glog.V(1).Infof("streamFromVolumeServers: %s/%s not found: %v", bucket, object, err)
s3err.WriteErrorResponse(w, r, s3err.ErrNoSuchKey)
return newStreamErrorWithResponse(err)
}
// serveObjectFromRemoteMount serves [offset, offset+size) straight from the
// mounted remote. served=false means nothing was written and the caller still
// owns the error path; once served is true the response is committed.
// owns the error path, with err saying why the origin could not be opened; once
// served is true the response is committed.
func (s3a *S3ApiServer) serveObjectFromRemoteMount(w http.ResponseWriter, r *http.Request, entry *filer_pb.Entry, bucket, object string, offset, size int64, isRangeRequest bool, totalSize int64, t0 time.Time) (served bool, err error) {
remoteReader, remoteErr := s3a.openRemoteStream(r.Context(), bucket, object, offset, size, nil)
if remoteErr != nil {
glog.Warningf("streamFromVolumeServers: origin stream %s/%s: %v", bucket, object, remoteErr)
return false, nil
return false, remoteErr
}
defer remoteReader.Close()
+42
View File
@@ -3,6 +3,7 @@ package s3api
import (
"bytes"
"context"
"errors"
"fmt"
"io"
"net"
@@ -610,6 +611,7 @@ type fakeStreamRemoteClient struct {
remote_storage.RemoteStorageClient
data []byte
stat *filer_pb.RemoteEntry // what the remote reports now; defaults to matching data
openErr error
gotLoc *remote_pb.RemoteStorageLocation
gotOffset int64
gotSize int64
@@ -624,6 +626,9 @@ func (c *fakeStreamRemoteClient) StatFile(loc *remote_pb.RemoteStorageLocation)
func (c *fakeStreamRemoteClient) ReadFileAsStream(ctx context.Context, loc *remote_pb.RemoteStorageLocation, offset int64, size int64) (io.ReadCloser, error) {
c.gotLoc, c.gotOffset, c.gotSize = loc, offset, size
if c.openErr != nil {
return nil, c.openErr
}
end := min(offset+size, int64(len(c.data)))
return io.NopCloser(bytes.NewReader(c.data[offset:end])), nil
}
@@ -827,3 +832,40 @@ func TestS3ColdReadStreamsFromOrigin(t *testing.T) {
}
})
}
// TestS3ReadOfObjectGoneFromRemote pins NoSuchKey for a remote-only entry whose
// remote object was deleted outside the filer, whether the filer's cache or the
// gateway's origin fallback learns it.
func TestS3ReadOfObjectGoneFromRemote(t *testing.T) {
entry := &filer_pb.Entry{
Name: "obj.bin",
Attributes: &filer_pb.FuseAttributes{FileSize: 10},
RemoteEntry: &filer_pb.RemoteEntry{RemoteSize: 10},
}
noWait := func(loc *remote_pb.RemoteStorageLocation) { loc.CacheWaitMs = proto.Int32(0) }
tests := []struct {
name string
cacheErr error
originErr error
mountOpts []func(*remote_pb.RemoteStorageLocation)
wantStatus int
}{
{"cache reports it gone", status.Error(codes.NotFound, "remote object not found"), remote_storage.ErrRemoteObjectNotFound, nil, http.StatusNotFound},
{"cache still filling, origin reports it gone", stillCachingErr, remote_storage.ErrRemoteObjectNotFound, nil, http.StatusNotFound},
{"zero cache wait, origin reports it gone", nil, remote_storage.ErrRemoteObjectNotFound, []func(*remote_pb.RemoteStorageLocation){noWait}, http.StatusNotFound},
{"origin denied", status.Error(codes.Internal, "assign: no free volumes"), errors.New("googleapi: Error 403: forbidden"), nil, http.StatusInternalServerError},
}
for i, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
prev := remote_storage.RemoteStorageClientMakers["faketest"]
remote_storage.RemoteStorageClientMakers["faketest"] = &fakeStreamRemoteMaker{client: &fakeStreamRemoteClient{openErr: tt.originErr}}
t.Cleanup(func() { remote_storage.RemoteStorageClientMakers["faketest"] = prev })
s3a := newRemoteCacheTestServer(startStreamThroughFiler(t, fmt.Sprintf("faketest-gone-%d", i), tt.cacheErr, tt.mountOpts...))
w := httptest.NewRecorder()
r := httptest.NewRequest(http.MethodGet, "/mybucket/dir/obj.bin", nil)
require.Error(t, s3a.streamFromVolumeServers(w, r, entry, "", "mybucket", "dir/obj.bin", ""))
assert.Equal(t, tt.wantStatus, w.Code)
})
}
}
+6
View File
@@ -468,6 +468,9 @@ func (fs *FilerServer) applyObjectMutation(ctx context.Context, m *filer_pb.Obje
case filer_pb.ObjectMutation_DELETE:
fullpath := util.NewFullPath(m.Directory, m.Name)
if m.KeepRemoteObject {
ctx = filer.WithKeepRemoteObject(ctx)
}
err := fs.filer.DeleteEntryMetaAndData(ctx, fullpath, m.IsRecursive, false, m.IsDeleteData, fromOtherCluster, signatures, 0)
if err != nil && err != filer_pb.ErrNotFound {
return err
@@ -896,6 +899,9 @@ func (fs *FilerServer) DeleteEntry(ctx context.Context, req *filer_pb.DeleteEntr
pathLock := fs.entryLockTable.AcquireLock("DeleteEntry", fullpath, util.ExclusiveLock)
defer fs.entryLockTable.ReleaseLock(fullpath, pathLock)
if req.KeepRemoteObject {
ctx = filer.WithKeepRemoteObject(ctx)
}
ctx, eventSink := filer.WithMetadataEventSink(ctx)
err = fs.filer.DeleteEntryMetaAndData(ctx, fullpath, req.IsRecursive, req.IgnoreRecursiveError, req.IsDeleteData, req.IsFromOtherCluster, req.Signatures, req.IfNotModifiedAfter)
resp = &filer_pb.DeleteEntryResponse{}
+81 -7
View File
@@ -5,6 +5,7 @@ import (
"errors"
"fmt"
"sort"
"strings"
"sync"
"time"
@@ -15,6 +16,9 @@ import (
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/remote_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
"github.com/seaweedfs/seaweedfs/weed/remote_storage"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_objectlock"
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
"github.com/seaweedfs/seaweedfs/weed/util"
"google.golang.org/grpc/codes"
@@ -50,12 +54,7 @@ func (fs *FilerServer) CacheRemoteObjectToLocalCluster(ctx context.Context, req
glog.V(2).Infof("CacheRemoteObjectToLocalCluster: shared result for %s", cacheKey)
}
if res.Err != nil {
// The sentinel would cross gRPC as codes.Unknown; make it canonical
// so remote callers can classify a vanished entry.
if errors.Is(res.Err, filer_pb.ErrNotFound) {
return nil, status.Error(codes.NotFound, res.Err.Error())
}
return nil, res.Err
return nil, cacheRemoteObjectError(res.Err)
}
if res.Val == nil {
return nil, fmt.Errorf("unexpected nil result from singleflight")
@@ -220,7 +219,7 @@ func (fs *FilerServer) doCacheRemoteObjectToLocalCluster(ctx context.Context, re
RemoteLocation: remoteLocation,
})
if fetchErr != nil {
return fmt.Errorf("volume server %s fetchAndWrite %s: %v", assignResult.Url, remoteLocation.Path, fetchErr)
return fetchAndWriteError(assignResult.Url, remoteLocation.Path, fetchErr)
}
etag = resp.ETag
return nil
@@ -277,6 +276,7 @@ func (fs *FilerServer) doCacheRemoteObjectToLocalCluster(ctx context.Context, re
fs.notePendingRemoteCacheVids(fileIds)
go fs.reclaimRemoteCacheSpace(fs.evictCtx(), entry.Remote.RemoteSize, nil)
}
fs.pruneEntryMissingFromRemote(ctx, lockPath, entry, err)
return nil, err
}
@@ -364,3 +364,77 @@ func (fs *FilerServer) resolveMountedRemote(ctx context.Context, dir, name strin
remoteLocation := filer.MapFullPathToRemoteStorageLocation(util.FullPath(localMountedDir), remoteStorageMountedLocation, util.FullPath(dir).Child(name))
return storageConf, remoteLocation, nil
}
// fetchAndWriteError restores the remote's not-found sentinel from the volume
// server's answer, so it is told apart from failures worth retrying.
func fetchAndWriteError(volumeServer, path string, err error) error {
if st, ok := status.FromError(err); ok && st.Code() == codes.NotFound && strings.Contains(st.Message(), remote_storage.ErrRemoteObjectNotFound.Error()) {
return fmt.Errorf("volume server %s fetchAndWrite %s: %w", volumeServer, path, remote_storage.ErrRemoteObjectNotFound)
}
return fmt.Errorf("volume server %s fetchAndWrite %s: %v", volumeServer, path, err)
}
// cacheRemoteObjectError makes the not-found sentinels canonical: they would
// cross gRPC as codes.Unknown, and remote callers classify a vanished entry or
// remote object by codes.NotFound.
func cacheRemoteObjectError(err error) error {
if errors.Is(err, filer_pb.ErrNotFound) || errors.Is(err, remote_storage.ErrRemoteObjectNotFound) {
return status.Error(codes.NotFound, err.Error())
}
return err
}
// pruneEntryMissingFromRemote converges the filer on the remote once the remote
// confirms the object is gone. It removes the entry only while it is still the
// one the failed fetch read: any write since then owns the path. The check and
// delete run as an object transaction so they land on the path's write owner,
// under the lock the S3 object writes take.
func (fs *FilerServer) pruneEntryMissingFromRemote(ctx context.Context, p util.FullPath, fetched *filer.Entry, fetchErr error) bool {
if !errors.Is(fetchErr, remote_storage.ErrRemoteObjectNotFound) || !isPrunableRemoteEntry(fetched) {
return false
}
dir, name := p.DirAndName()
resp, err := fs.ObjectTransaction(ctx, &filer_pb.ObjectTransactionRequest{
LockKey: string(p),
RouteKey: entryRouteKey(p),
Condition: &filer_pb.WriteCondition{Clauses: []*filer_pb.WriteCondition_Clause{{
Kind: filer_pb.WriteCondition_IF_ENTRY_EQUAL,
ExpectedEntry: fetched.ToProtoEntry(),
}}},
Mutations: []*filer_pb.ObjectMutation{{
Type: filer_pb.ObjectMutation_DELETE,
Directory: dir,
Name: name,
KeepRemoteObject: true,
}},
})
if err != nil {
glog.WarningfCtx(ctx, "prune %s missing from remote: %v", p, err)
return false
}
if resp.Error != "" {
if resp.ErrorCode != filer_pb.FilerError_PRECONDITION_FAILED {
glog.WarningfCtx(ctx, "prune %s missing from remote: %s", p, resp.Error)
}
return false
}
glog.V(0).InfofCtx(ctx, "pruned %s: its remote object no longer exists", p)
return true
}
// isPrunableRemoteEntry admits only a file whose content lives solely on the
// remote. Version entries are left alone because their parent's latest-version
// pointer would dangle, and locked objects because pruning would bypass the
// lock.
func isPrunableRemoteEntry(entry *filer.Entry) bool {
if entry == nil || entry.IsDirectory() || !entry.IsInRemoteOnly() {
return false
}
if len(entry.Content) > 0 || len(entry.HardLinkId) > 0 {
return false
}
if dir, _ := entry.FullPath.DirAndName(); strings.HasSuffix(dir, s3_constants.VersionsFolder) {
return false
}
return !s3_objectlock.EntryHasActiveLock(entry.ToProtoEntry(), time.Now())
}
@@ -0,0 +1,229 @@
package weed_server
import (
"context"
"errors"
"fmt"
"os"
"strconv"
"sync"
"testing"
"time"
"github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager"
"github.com/seaweedfs/seaweedfs/weed/filer"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/remote_pb"
"github.com/seaweedfs/seaweedfs/weed/remote_storage"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
"github.com/seaweedfs/seaweedfs/weed/util"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"
)
func TestRemoteNotFoundCrossesTheVolumeServerHop(t *testing.T) {
loc := &remote_pb.RemoteStorageLocation{Name: "gcs1", Bucket: "bucket", Path: "/obj.bin"}
tests := map[string]struct {
volumeErr error
notFound bool
}{
"remote object gone": {volumeRemoteReadError(loc, fmt.Errorf("read: %w", remote_storage.ErrRemoteObjectNotFound)), true},
"remote read failure": {volumeRemoteReadError(loc, errors.New("googleapi: Error 403: forbidden")), false},
"not found for other uses": {status.Error(codes.NotFound, "volume not found"), false},
"deadline": {status.Error(codes.DeadlineExceeded, "context deadline exceeded"), false},
"unavailable": {status.Error(codes.Unavailable, "connection refused"), false},
}
for name, tt := range tests {
t.Run(name, func(t *testing.T) {
filerErr := fetchAndWriteError("volume:8080", loc.Path, tt.volumeErr)
assert.Equal(t, tt.notFound, errors.Is(filerErr, remote_storage.ErrRemoteObjectNotFound))
assert.Equal(t, tt.notFound, status.Code(cacheRemoteObjectError(filerErr)) == codes.NotFound)
})
}
}
// mountTestRemoteClient is a remote where every object is gone; it records
// every DeleteFile.
type mountTestRemoteClient struct {
remote_storage.RemoteStorageClient
mu sync.Mutex
deleted []string
}
func (c *mountTestRemoteClient) StatFile(*remote_pb.RemoteStorageLocation) (*filer_pb.RemoteEntry, error) {
return nil, remote_storage.ErrRemoteObjectNotFound
}
func (c *mountTestRemoteClient) DeleteFile(loc *remote_pb.RemoteStorageLocation) error {
c.mu.Lock()
defer c.mu.Unlock()
c.deleted = append(c.deleted, loc.Path)
return nil
}
func (c *mountTestRemoteClient) deletedPaths() []string {
c.mu.Lock()
defer c.mu.Unlock()
return append([]string(nil), c.deleted...)
}
// newRemoteMountTestServer mounts /buckets/b on a fresh mountTestRemoteClient.
func newRemoteMountTestServer(t *testing.T) (*FilerServer, *renameTestStore, *mountTestRemoteClient) {
t.Helper()
store := newRenameTestStore()
f := newRenameTestFiler(t, store)
client := &mountTestRemoteClient{}
f.BuildGuardedRemoteClient = func(context.Context, *remote_pb.RemoteConf, bool) (remote_storage.RemoteStorageClient, error) {
return client, nil
}
mapping, err := proto.Marshal(&remote_pb.RemoteStorageMapping{Mappings: map[string]*remote_pb.RemoteStorageLocation{
"/buckets/b": {Name: "r1", Bucket: "origin", Path: "/"},
}})
require.NoError(t, err)
conf, err := proto.Marshal(&remote_pb.RemoteConf{Name: "r1", Type: "s3"})
require.NoError(t, err)
for i, c := range map[string][]byte{filer.REMOTE_STORAGE_MOUNT_FILE: mapping, "r1" + filer.REMOTE_STORAGE_CONF_SUFFIX: conf} {
entry := newFileEntry(filer.DirectoryEtcRemote+"/"+i, 1)
entry.Content = c
store.entries[string(entry.FullPath)] = entry
}
require.NoError(t, f.RemoteStorage.LoadRemoteStorageConfigurationsAndMapping(f))
return &FilerServer{filer: f, option: &FilerOption{}, entryLockTable: util.NewLockTable[util.FullPath]()}, store, client
}
func remoteOnlyFileEntry(path string) *filer.Entry {
entry := newFileEntry(path, 10)
entry.FileSize = 10
entry.Remote = &filer_pb.RemoteEntry{RemoteSize: 10, RemoteMtime: 1700000000}
dir, _ := entry.FullPath.DirAndName()
return filer.FromPbEntry(dir, entry.ToProtoEntry())
}
func TestPruneEntryMissingFromRemote(t *testing.T) {
gone := fmt.Errorf("volume server v fetchAndWrite /obj.bin: %w", remote_storage.ErrRemoteObjectNotFound)
retainUntil := func(d time.Duration) []byte { return []byte(strconv.FormatInt(time.Now().Add(d).Unix(), 10)) }
tests := []struct {
name string
fetchErr error
mutate func(stored *filer.Entry)
pruned bool
}{
{"remote confirms the object gone", gone, nil, true},
{"expired retention", gone, func(e *filer.Entry) {
e.Extended = map[string][]byte{s3_constants.ExtObjectLockModeKey: []byte(s3_constants.RetentionModeGovernance), s3_constants.ExtRetentionUntilDateKey: retainUntil(-time.Hour)}
}, true},
{"deadline", context.DeadlineExceeded, nil, false},
{"unavailable", fmt.Errorf("fetchAndWrite: %v", status.Error(codes.Unavailable, "connection refused")), nil, false},
{"permission denied", errors.New("googleapi: Error 403: forbidden"), nil, false},
{"entry vanished", filer_pb.ErrNotFound, nil, false},
{"cached locally", gone, func(e *filer.Entry) { e.Chunks = []*filer_pb.FileChunk{{FileId: "3,01637037d6", Size: 10}} }, false},
{"inline content", gone, func(e *filer.Entry) { e.Content = []byte("0123456789") }, false},
{"hard link", gone, func(e *filer.Entry) { e.HardLinkId = filer.HardLinkId([]byte("link")) }, false},
{"directory", gone, func(e *filer.Entry) { e.Mode |= os.ModeDir }, false},
{"object version", gone, func(e *filer.Entry) {
e.FullPath = util.FullPath("/buckets/b/obj.bin" + s3_constants.VersionsFolder + "/v_1")
}, false},
{"active retention", gone, func(e *filer.Entry) {
e.Extended = map[string][]byte{s3_constants.ExtObjectLockModeKey: []byte(s3_constants.RetentionModeGovernance), s3_constants.ExtRetentionUntilDateKey: retainUntil(time.Hour)}
}, false},
{"legal hold", gone, func(e *filer.Entry) {
e.Extended = map[string][]byte{s3_constants.ExtLegalHoldKey: []byte(s3_constants.LegalHoldOn)}
}, false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
fs, store, client := newRemoteMountTestServer(t)
entry := remoteOnlyFileEntry("/buckets/b/obj.bin")
if tt.mutate != nil {
tt.mutate(entry)
}
store.entries[string(entry.FullPath)] = entry.ShallowClone()
assert.Equal(t, tt.pruned, fs.pruneEntryMissingFromRemote(context.Background(), entry.FullPath, entry, tt.fetchErr))
_, err := fs.filer.FindEntry(context.Background(), entry.FullPath)
assert.Equal(t, tt.pruned, errors.Is(err, filer_pb.ErrNotFound), "find: %v", err)
assert.Empty(t, client.deletedPaths())
})
}
t.Run("entry rewritten since the fetch", func(t *testing.T) {
fs, store, _ := newRemoteMountTestServer(t)
fetched := remoteOnlyFileEntry("/buckets/b/obj.bin")
rewritten := fetched.ShallowClone()
rewritten.Remote = nil
rewritten.Chunks = []*filer_pb.FileChunk{{FileId: "3,01637037d6", Size: 10}}
store.entries[string(fetched.FullPath)] = rewritten
assert.False(t, fs.pruneEntryMissingFromRemote(context.Background(), fetched.FullPath, fetched, gone))
})
}
// TestPruneLeavesTheRemoteAlone also deletes an entry the ordinary way, so a
// harness that never reached the remote could not pass it.
func TestPruneLeavesTheRemoteAlone(t *testing.T) {
fs, store, client := newRemoteMountTestServer(t)
pruned := remoteOnlyFileEntry("/buckets/b/pruned.bin")
deleted := remoteOnlyFileEntry("/buckets/b/deleted.bin")
store.entries[string(pruned.FullPath)] = pruned.ShallowClone()
store.entries[string(deleted.FullPath)] = deleted.ShallowClone()
queue := &captureQueue{}
swapNotificationQueue(t, queue)
require.True(t, fs.pruneEntryMissingFromRemote(context.Background(), pruned.FullPath, pruned, remote_storage.ErrRemoteObjectNotFound))
events := queue.snapshot()
require.Len(t, events, 1)
assert.False(t, events[0].notification.IsFromOtherCluster, "filer.sync must replicate the prune")
assert.Contains(t, events[0].notification.OldEntry.Extended, filer.ExtKeepRemoteObjectKey, "filer.remote.sync must not delete the remote object")
require.NoError(t, fs.filer.DeleteEntryMetaAndData(context.Background(), deleted.FullPath, false, false, false, false, nil, 0))
assert.Equal(t, []string{"/deleted.bin"}, client.deletedPaths())
}
func TestReplicatedMetadataOnlyDeleteStaysMarked(t *testing.T) {
for _, keep := range []bool{true, false} {
t.Run(fmt.Sprintf("keep_remote_object=%v", keep), func(t *testing.T) {
fs, store, client := newRemoteMountTestServer(t)
entry := remoteOnlyFileEntry("/buckets/b/obj.bin")
store.entries[string(entry.FullPath)] = entry.ShallowClone()
resp, err := fs.DeleteEntry(context.Background(), &filer_pb.DeleteEntryRequest{
Directory: "/buckets/b",
Name: "obj.bin",
IsFromOtherCluster: true,
KeepRemoteObject: keep,
})
require.NoError(t, err)
require.Empty(t, resp.Error)
assert.Equal(t, keep, filer.IsMetadataOnlyDelete(resp.MetadataEvent.EventNotification.OldEntry))
assert.Empty(t, client.deletedPaths())
})
}
}
func TestPruneOnANonOwnerLeavesTheDeleteToTheOwner(t *testing.T) {
fs, store, _ := newRemoteMountTestServer(t)
const self = pb.ServerAddress("127.0.0.1:18888")
fs.option.Host = self
fs.grpcDialOption = grpc.WithTransportCredentials(insecure.NewCredentials())
fs.filer.Dlm = lock_manager.NewDistributedLockManager(self)
fs.filer.Dlm.LockRing.SetSnapshot([]pb.ServerAddress{self, "127.0.0.1:18889"}, 1)
var entry *filer.Entry
for i := 0; i < 4000 && entry == nil; i++ {
p := util.FullPath(fmt.Sprintf("/buckets/b/obj-%d.bin", i))
if fs.filer.Dlm.LockRing.WriteOwner(entryRouteKey(p)) != self {
entry = remoteOnlyFileEntry(string(p))
}
}
require.NotNil(t, entry, "no path owned by the peer")
store.entries[string(entry.FullPath)] = entry.ShallowClone()
assert.False(t, fs.pruneEntryMissingFromRemote(context.Background(), entry.FullPath, entry, remote_storage.ErrRemoteObjectNotFound))
assert.Contains(t, store.entries, string(entry.FullPath), "only the write owner may delete the entry")
}
+14 -1
View File
@@ -14,6 +14,8 @@ import (
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"github.com/seaweedfs/seaweedfs/weed/operation"
"github.com/seaweedfs/seaweedfs/weed/pb/remote_pb"
@@ -525,7 +527,7 @@ func (vs *VolumeServer) FetchAndWriteNeedle(ctx context.Context, req *volume_ser
data, readRemoteErr = client.ReadFile(remoteStorageLocation, req.Offset, req.Size)
}
if readRemoteErr != nil {
return nil, fmt.Errorf("read from remote %+v: %w", remoteStorageLocation, readRemoteErr)
return nil, volumeRemoteReadError(remoteStorageLocation, readRemoteErr)
}
// The chunk is recorded with the requested size, so a short read would be
// cached as a full-size chunk with a zero-padded or truncated tail. Fail
@@ -617,3 +619,14 @@ func (vs *VolumeServer) FetchAndWriteNeedle(ctx context.Context, req *volume_ser
return resp, err
}
// volumeRemoteReadError keeps a confirmed missing object distinguishable across
// gRPC, where a wrapped sentinel would arrive as codes.Unknown text. The
// sentinel's text in the message is what the filer matches, so a NotFound for
// any other reason cannot pass for it.
func volumeRemoteReadError(loc *remote_pb.RemoteStorageLocation, err error) error {
if errors.Is(err, remote_storage.ErrRemoteObjectNotFound) {
return status.Errorf(codes.NotFound, "read from remote %+v: %v", loc, remote_storage.ErrRemoteObjectNotFound)
}
return fmt.Errorf("read from remote %+v: %w", loc, err)
}