Files
seaweedfs/weed/server/filer_grpc_server_remote.go
T
Chris LuGitHubDevin <158243242+devin-ai-integration[bot]@users.noreply.github.com>Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
757917f564 filer: evict remote-cached objects under storage pressure (#11515)
* filer: identify remote-mounted entries safe to drop under disk pressure

ListEvictableRemoteEntries walks every mounted directory directly on the
filer store (no lazy remote listing) and returns entries that hold local
chunks fully synchronized with remote, ordered oldest-cached first.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: evict remote-cached chunks oldest-first and vacuum the garbage

uncacheRemoteEntry applies the same transition remote.uncache does -
cleared chunks plus a reset LastLocalSyncTsNs under the entry path lock -
and evictRemoteCachedEntries serializes passes over all mounts until a
byte target is met. Aged victims are preferred; a second pass accepts any
synchronized cached entry when aged ones cannot cover the request, since
a failed read is worse than a dropped hot object.

Cleared chunks only become disk space after compaction, so
reclaimRemoteCacheSpace pairs each pass with a rate-limited VacuumVolume
call that also picks up orphaned partial fills.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: trigger remote cache eviction under storage pressure

A periodic check (30s) reads disk usage from master topology and evicts
remote-mounted cached chunks once any disk crosses
-filer.remoteCacheEvictThreshold (default 0.9; 0 disables), with a vacuum
pass to reclaim the tombstoned needles.

The cold-read cache path also kicks the same reclaim when a fill fails on
exhausted volumes - the request still falls back to streaming from the
remote, but the cache stops being permanently wedged full.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: flush deletion queue before remote cache vacuum

Vacuum ran immediately after eviction while evicted file IDs still sat
in the asynchronous deletion queue, so compaction saw no garbage and the
cache stayed wedged. Flush the queue synchronously first and shorten the
vacuum cooldown so sustained pressure does not wait five minutes between
reclaim passes.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test: cover remote cache eviction under capacity pressure

Unit tests pin the eligibility filter and oldest-first ordering; the
integration test runs a constrained two-node setup that saturates the
cache, verifies the oldest synced entry is evicted and vacuumed, and
that a later read re-caches it.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: coalesce remote cache reclaim passes

A failed cache fill used to queue behind any in-flight eviction,
stacking full mount traversals during a write-failure storm. Skip the
pass when one is already running; the caller falls back to streaming
from remote regardless.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: stop the remote cache janitor on shutdown

The eviction ticker kept running after Shutdown closed the metadata
store and could traverse a closed store. Give the janitor a context
cancelled from Shutdown and propagate it into its master RPCs and
traversals.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: vacuum only tombstoned volumes and retry deferred passes

VacuumVolume with no volume id swept every collection, compacting
volumes unrelated to the cache fill that failed. Now the reclaim path
collects the vids of file ids actually flushed from the deletion queue
and compacts only those. Vids that land inside the vacuum cooldown stay
in a pending set the janitor retries on each tick, so chunks evicted
just after a sweep are not stranded until the next pressure event.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: count only pressured disks when evicting remote cache

The janitor measured the largest excess on one disk but let bytes on
healthy disks satisfy the reclaim target. Split the topology disk view
per physical disk and count only chunk bytes whose volumes sit on an
over-threshold disk; entries contributing nothing there are skipped.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: compare remote cache sync time at nanosecond precision

Second-precision mtime comparisons let a local write in the same second
as the last sync still qualify as evictable, discarding unsynced
changes. Compare LastLocalSyncTsNs against full-precision mtime
(mtime_ns round-trips through the entry codec), and apply the same fix
to remote.uncache's inline check.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: invalidate remote sync stamp on local content change

A local overwrite that keeps the remote entry's LastLocalSyncTsNs looks
evictable even though the remote copy no longer matches, and some write
paths stamp mtime at second precision so a timestamp comparison cannot
catch it. UpdateEntry now clears the stamp when chunks change without a
fresh stamp, leaving replicated updates authoritative.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: bound remote cache master rpcs and vacuum all evicted garbage

VolumeList and VacuumVolume now run under a 30s context so a stalled
master cannot wedge the eviction janitor. The targeted vacuum drops the
garbage threshold so volumes with under 10% deleted bytes still compact.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test: tolerate straggler fills in remote cache eviction test

Detached fills from the concurrent wave keep racing the final checks:
live chunks legitimately fill both volumes, and a re-cached object can
be evicted again before its commit is observed.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: start remote cache eviction loop after filer init

The janitor's first tick dereferences fs.filer; starting the goroutine
before NewFiler assigns it could panic when startup exceeds an interval.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: keep remote cache vacuum intent across retries

Evicted entries now record their chunk volumes for vacuum directly, so
the intent survives whoever consumes the shared deletion queue first.
A pending volume keeps several vacuum attempts so tombstones that land
late are still compacted, and the janitor retries pending volumes under
the reclaim mutex instead of flushing unrelated deletes every tick.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: bound each remote cache vacuum request independently

A shared 30s deadline across pending volumes let one slow compaction
cancel the rest. Each VacuumVolume now gets its own context, and pending
volumes keep more attempts since the master reports request acceptance
rather than compaction.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* test: tighten remote cache reclamation bound

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: treat chunk timestamp changes as content changes

chunksEqual now also compares ModifiedTsNs so an update that rewrites a
chunk record still invalidates the remote sync stamp.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: run remote cache queue flush under the reclaim context

BatchDelete for flushed file ids now uses the caller's context instead of
context.Background(), so a reclaim pass bounded by shutdown or timeout
stops its deletes too. Other callers keep their existing behavior.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: scope remote cache vacuum to evicted volumes

The flush no longer feeds the shared deletion queue's ids into the
pending set — only evicted chunks' volumes are tracked, so ordinary
deletions no longer pick up repeated vacuum attempts. The flush also
runs under a shutdown-immune bounded context and is skipped when no
volume is pending.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: retry remote cache vacuums even after unmount

Pending volumes were only retried while a remote mount existed; removing
the last mount skipped every later pass and left evicted bytes allocated.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

* filer: reclaim partial cache fills that run out of capacity

A fill that fails midway queues its written chunks for deletion, but
when no entries remain evictable the reclaim pass found no pending
volumes and skipped the flush and vacuum entirely, leaving the partial
garbage to the slow periodic vacuum while the disk stayed full. Mark
the failed fill's chunk volumes pending so the pass tombstones and
compacts them even when nothing was evicted.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>

---------

Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
2026-09-29 22:05:42 +08:00

367 lines
13 KiB
Go

package weed_server
import (
"context"
"errors"
"fmt"
"sort"
"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/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 {
// 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
}
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 fmt.Errorf("volume server %s fetchAndWrite %s: %v", 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)
}
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
}