Files
seaweedfs/weed/server/filer_grpc_server_remote.go
T
Peter Dodd 5f1a742726 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.
2026-10-10 11:00:04 +08:00

441 lines
16 KiB
Go

package weed_server
import (
"context"
"errors"
"fmt"
"sort"
"strings"
"sync"
"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"
"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"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/proto"
)
func (fs *FilerServer) CacheRemoteObjectToLocalCluster(ctx context.Context, req *filer_pb.CacheRemoteObjectToLocalClusterRequest) (*filer_pb.CacheRemoteObjectToLocalClusterResponse, error) {
// Use singleflight to deduplicate concurrent caching requests for the same object.
// This benefits all clients: S3 API, filer HTTP, Hadoop, etc.
cacheKey := req.Directory + "/" + req.Name
// Detach from caller ctx: on failure the error path deletes every chunk
// already written, so cancelling mid-download loses all progress. For
// blobs large enough that the download outlasts the caller's timeout
// the retry loop never converges.
bgCtx := context.WithoutCancel(ctx)
// DoChan (vs Do) so the caller can bail out on ctx.Done() while the
// singleflight goroutine keeps caching on bgCtx; otherwise this handler
// goroutine stays blocked for the full download after the client is gone.
ch := fs.remoteCacheGroup.DoChan(cacheKey, func() (interface{}, error) {
return fs.doCacheRemoteObjectToLocalCluster(bgCtx, req)
})
select {
case <-ctx.Done():
// Caller gave up; the detached cache keeps running and a later
// request will find the entry cached (or join the same singleflight).
return nil, ctx.Err()
case res := <-ch:
if res.Shared {
glog.V(2).Infof("CacheRemoteObjectToLocalCluster: shared result for %s", cacheKey)
}
if res.Err != nil {
return nil, cacheRemoteObjectError(res.Err)
}
if res.Val == nil {
return nil, fmt.Errorf("unexpected nil result from singleflight")
}
resp, ok := res.Val.(*filer_pb.CacheRemoteObjectToLocalClusterResponse)
if !ok {
return nil, fmt.Errorf("unexpected result type from singleflight")
}
return resp, nil
}
}
// doCacheRemoteObjectToLocalCluster performs the actual caching operation.
// This is called from singleflight, so only one instance runs per object.
func (fs *FilerServer) doCacheRemoteObjectToLocalCluster(ctx context.Context, req *filer_pb.CacheRemoteObjectToLocalClusterRequest) (*filer_pb.CacheRemoteObjectToLocalClusterResponse, error) {
lockPath := util.JoinPath(req.Directory, req.Name)
entry, logTsNs, err := fs.fencedFindEntry(ctx, lockPath)
if err == filer_pb.ErrNotFound {
return nil, err
}
if err != nil {
return nil, fmt.Errorf("find entry %s/%s: %v", req.Directory, req.Name, err)
}
resp := &filer_pb.CacheRemoteObjectToLocalClusterResponse{LogTsNs: logTsNs, LogSignature: fs.filer.Signature}
// Early return if not a remote-only object or already cached
if entry.Remote == nil || entry.Remote.RemoteSize == 0 {
resp.Entry = entry.ToProtoEntry()
return resp, nil
}
if len(entry.GetChunks()) > 0 {
// Already has local chunks - already cached
glog.V(2).Infof("CacheRemoteObjectToLocalCluster: %s/%s already cached (%d chunks)", req.Directory, req.Name, len(entry.GetChunks()))
resp.Entry = entry.ToProtoEntry()
return resp, nil
}
glog.V(1).Infof("CacheRemoteObjectToLocalCluster: caching %s/%s (remote size: %d)", req.Directory, req.Name, entry.Remote.RemoteSize)
storageConf, remoteLocation, err := fs.resolveMountedRemote(ctx, req.Directory, req.Name)
if err != nil {
return nil, err
}
// detect storage option
so, err := fs.detectStorageOption(ctx, req.Directory, "", "", 0, "", "", "", "")
if err != nil {
return resp, err
}
assignRequest, altRequest := so.ToAssignRequests(1)
// adaptive chunk size: target ~32 chunks per file to balance
// per-chunk overhead (volume assign, gRPC, needle write) against parallelism
chunkSize := int64(5 * 1024 * 1024) // 5MB floor
maxChunkSize := int64(fs.option.MaxMB) * 1024 * 1024
if maxChunkSize < chunkSize {
maxChunkSize = chunkSize
}
targetChunks := int64(32)
if entry.Remote.RemoteSize/targetChunks > chunkSize {
chunkSize = entry.Remote.RemoteSize / targetChunks
if chunkSize > maxChunkSize {
chunkSize = maxChunkSize
}
}
// final safety check: ensure no more than 1000 chunks
if (entry.Remote.RemoteSize+chunkSize-1)/chunkSize > 1000 {
chunkSize = (entry.Remote.RemoteSize + 999) / 1000
}
// Now that chunkSize is known, hint it to the master so per-chunk
// assigns don't fall back to the 1 MB default estimate. Slightly over-
// estimates for the final partial chunk (< chunkSize) by design.
assignRequest.ExpectedDataSize = uint64(chunkSize)
if altRequest != nil {
altRequest.ExpectedDataSize = uint64(chunkSize)
}
var chunks []*filer_pb.FileChunk
var chunksMu sync.Mutex
var fetchAndWriteErr error
var wg sync.WaitGroup
chunkConcurrency := int(req.ChunkConcurrency)
if chunkConcurrency <= 0 {
chunkConcurrency = 8
} else if chunkConcurrency > 1024 {
glog.V(0).Infof("capping chunkConcurrency from %d to 1024", chunkConcurrency)
chunkConcurrency = 1024
}
downloadConcurrency := req.DownloadConcurrency
if downloadConcurrency > 1024 {
glog.V(0).Infof("capping downloadConcurrency from %d to 1024", downloadConcurrency)
downloadConcurrency = 1024
}
limitedConcurrentExecutor := util.NewLimitedConcurrentExecutor(chunkConcurrency)
for offset := int64(0); offset < entry.Remote.RemoteSize; offset += chunkSize {
localOffset := offset
wg.Add(1)
limitedConcurrentExecutor.Execute(func() {
defer wg.Done()
size := chunkSize
if localOffset+chunkSize > entry.Remote.RemoteSize {
size = entry.Remote.RemoteSize - localOffset
}
// assign one volume server
assignResult, err := operation.Assign(ctx, fs.filer.GetMaster, fs.grpcDialOption, assignRequest, altRequest)
if err != nil {
chunksMu.Lock()
if fetchAndWriteErr == nil {
fetchAndWriteErr = err
}
chunksMu.Unlock()
return
}
if assignResult.Error != "" {
chunksMu.Lock()
if fetchAndWriteErr == nil {
fetchAndWriteErr = fmt.Errorf("assign: %v", assignResult.Error)
}
chunksMu.Unlock()
return
}
fileId, parseErr := needle.ParseFileIdFromString(assignResult.Fid)
if parseErr != nil {
chunksMu.Lock()
if fetchAndWriteErr == nil {
fetchAndWriteErr = fmt.Errorf("unrecognized file id %s: %v", assignResult.Fid, parseErr)
}
chunksMu.Unlock()
return
}
var replicas []*volume_server_pb.FetchAndWriteNeedleRequest_Replica
for _, r := range assignResult.Replicas {
replicas = append(replicas, &volume_server_pb.FetchAndWriteNeedleRequest_Replica{
Url: r.Url,
PublicUrl: r.PublicUrl,
GrpcPort: int32(r.GrpcPort),
})
}
// tell filer to tell volume server to download into needles
assignedServerAddress := pb.NewServerAddressWithGrpcPort(assignResult.Url, assignResult.GrpcPort)
var etag string
err = operation.WithVolumeServerClient(false, assignedServerAddress, fs.grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
resp, fetchErr := volumeServerClient.FetchAndWriteNeedle(context.Background(), &volume_server_pb.FetchAndWriteNeedleRequest{
VolumeId: uint32(fileId.VolumeId),
NeedleId: uint64(fileId.Key),
Cookie: uint32(fileId.Cookie),
Offset: localOffset,
Size: size,
Replicas: replicas,
Auth: string(assignResult.Auth),
DownloadConcurrency: downloadConcurrency,
RemoteConf: storageConf,
RemoteLocation: remoteLocation,
})
if fetchErr != nil {
return fetchAndWriteError(assignResult.Url, remoteLocation.Path, fetchErr)
}
etag = resp.ETag
return nil
})
if err != nil {
chunksMu.Lock()
if fetchAndWriteErr == nil {
fetchAndWriteErr = err
}
chunksMu.Unlock()
return
}
chunk := &filer_pb.FileChunk{
FileId: assignResult.Fid,
Offset: localOffset,
Size: uint64(size),
ModifiedTsNs: time.Now().UnixNano(),
ETag: etag,
Fid: &filer_pb.FileId{
VolumeId: uint32(fileId.VolumeId),
FileKey: uint64(fileId.Key),
Cookie: uint32(fileId.Cookie),
},
}
chunksMu.Lock()
chunks = append(chunks, chunk)
chunksMu.Unlock()
})
}
wg.Wait()
chunksMu.Lock()
err = fetchAndWriteErr
// Sort chunks by offset to maintain file order
sort.Slice(chunks, func(i, j int) bool {
return chunks[i].Offset < chunks[j].Offset
})
chunksMu.Unlock()
if err != nil {
// Clean up any chunks that were successfully written before the error.
// Without this, partial downloads leave orphaned needles in volume servers
// that accumulate across retry cycles and cannot be reclaimed by vacuum.
if len(chunks) > 0 {
fs.filer.DeleteUncommittedChunks(ctx, chunks)
}
if fs.option.RemoteCacheEvictThreshold > 0 && isRemoteCacheCapacityError(err) {
fileIds := make([]string, 0, len(chunks))
for _, chunk := range chunks {
fileIds = append(fileIds, chunk.GetFileIdString())
}
fs.notePendingRemoteCacheVids(fileIds)
go fs.reclaimRemoteCacheSpace(fs.evictCtx(), entry.Remote.RemoteSize, nil)
}
fs.pruneEntryMissingFromRemote(ctx, lockPath, entry, err)
return nil, err
}
// Commit under the mutation path lock so a fenced lookup cannot land
// between the store update and its notification, handing out
// under-versioned state. Re-read under it: the entry may have changed
// during the unlocked download, and the stale base would clobber it.
commitLock := fs.entryLockTable.AcquireLock("CacheRemoteObjectToLocalCluster", lockPath, util.ExclusiveLock)
defer fs.entryLockTable.ReleaseLock(lockPath, commitLock)
commitLogTsNs := time.Now().UnixNano()
current, err := fs.filer.FindEntry(ctx, lockPath)
if err != nil {
fs.filer.DeleteUncommittedChunks(ctx, chunks)
if err == filer_pb.ErrNotFound {
// Deleted while the download ran; keep the sentinel so callers
// still surface a 404 rather than a generic failure.
return nil, err
}
return nil, fmt.Errorf("find entry %s before commit: %v", lockPath, err)
}
if !filer.EqualEntry(current, entry) {
// Changed during the download: that writer supersedes the cached
// content. Return the current state, fenced at this read.
fs.filer.DeleteUncommittedChunks(ctx, chunks)
resp.Entry = current.ToProtoEntry()
resp.LogTsNs = commitLogTsNs
resp.LogSignature = fs.filer.Signature
return resp, nil
}
garbage := entry.GetChunks()
newEntry := entry.ShallowClone()
newEntry.Chunks = chunks
newEntry.Remote = proto.Clone(entry.Remote).(*filer_pb.RemoteEntry)
newEntry.Remote.LastLocalSyncTsNs = time.Now().UnixNano()
// this skips meta data log events
if err := fs.filer.Store.UpdateEntry(context.Background(), newEntry); err != nil {
fs.filer.DeleteUncommittedChunks(ctx, chunks)
return nil, err
}
fs.filer.DeleteChunks(ctx, entry.FullPath, garbage)
ctx, eventSink := filer.WithMetadataEventSink(ctx)
fs.filer.NotifyUpdateEvent(ctx, entry, newEntry, true, false, nil)
resp.Entry = newEntry.ToProtoEntry()
resp.MetadataEvent = eventSink.Last()
resp.LogTsNs = commitLogTsNs
resp.LogSignature = fs.filer.Signature
return resp, nil
}
// resolveMountedRemote reads /etc/remote fresh (so conf changes need no restart)
// and maps dir/name to its remote storage conf and remote location.
func (fs *FilerServer) resolveMountedRemote(ctx context.Context, dir, name string) (*remote_pb.RemoteConf, *remote_pb.RemoteStorageLocation, error) {
mappingEntry, err := fs.filer.FindEntry(ctx, util.JoinPath(filer.DirectoryEtcRemote, filer.REMOTE_STORAGE_MOUNT_FILE))
if err != nil {
return nil, nil, err
}
mappings, err := filer.UnmarshalRemoteStorageMappings(mappingEntry.Content)
if err != nil {
return nil, nil, err
}
localMountedDir, remoteStorageMountedLocation, err := filer.FindMountedRemoteMapping(mappings, dir)
if err != nil {
return nil, nil, err
}
storageConfEntry, err := fs.filer.FindEntry(ctx, util.JoinPath(filer.DirectoryEtcRemote, remoteStorageMountedLocation.Name+filer.REMOTE_STORAGE_CONF_SUFFIX))
if err != nil {
return nil, nil, err
}
storageConf := &remote_pb.RemoteConf{}
if unMarshalErr := proto.Unmarshal(storageConfEntry.Content, storageConf); unMarshalErr != nil {
return nil, nil, fmt.Errorf("unmarshal remote storage conf %s/%s: %v", filer.DirectoryEtcRemote, remoteStorageMountedLocation.Name+filer.REMOTE_STORAGE_CONF_SUFFIX, unMarshalErr)
}
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())
}