From 757917f5640f277eee5afbf25339ca4f5806391f Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Tue, 29 Sep 2026 22:05:42 +0800 Subject: [PATCH] filer: evict remote-cached objects under storage pressure (#11515) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * 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> --- .../remote_cache/remote_cache_evict_test.go | 231 ++++++++++++++ weed/command/filer.go | 3 + weed/command/mini.go | 1 + weed/command/server.go | 1 + weed/filer/filer.go | 18 +- weed/filer/filer_deletion.go | 35 +- weed/filer/filer_inode_test.go | 6 +- weed/filer/filer_remote_evict.go | 89 ++++++ weed/filer/filer_remote_evict_test.go | 164 ++++++++++ weed/operation/delete_content.go | 10 +- weed/s3api/s3api_object_handlers_put.go | 2 +- weed/server/filer_grpc_server.go | 8 +- weed/server/filer_grpc_server_remote.go | 8 + weed/server/filer_server.go | 17 + weed/server/filer_server_remote_evict.go | 300 ++++++++++++++++++ weed/shell/command_fs_merge_volumes.go | 2 +- weed/shell/command_remote_uncache.go | 2 +- weed/shell/command_volume_check_disk.go | 1 + weed/shell/command_volume_fsck.go | 2 +- weed/stats/metrics.go | 9 + 20 files changed, 875 insertions(+), 34 deletions(-) create mode 100644 test/s3/remote_cache/remote_cache_evict_test.go create mode 100644 weed/filer/filer_remote_evict.go create mode 100644 weed/filer/filer_remote_evict_test.go create mode 100644 weed/server/filer_server_remote_evict.go diff --git a/test/s3/remote_cache/remote_cache_evict_test.go b/test/s3/remote_cache/remote_cache_evict_test.go new file mode 100644 index 000000000..49060319f --- /dev/null +++ b/test/s3/remote_cache/remote_cache_evict_test.go @@ -0,0 +1,231 @@ +package remote_cache + +import ( + "bytes" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "os/exec" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + "github.com/aws/aws-sdk-go/aws" + "github.com/aws/aws-sdk-go/service/s3" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +const ( + evictRemoteS3 = "http://localhost:28334" + evictRemoteMaster = "30334" + evictRemoteFiler = "28889" + evictRemoteVolume = "30341" + evictRemoteWebdav = "27334" + evictRemoteMetrics = "30326" + + evictPrimaryS3 = "http://localhost:28333" + evictPrimaryMaster = "30333" + evictPrimaryFiler = "28887" + evictPrimaryVolume = "30340" + evictPrimaryWebdav = "27333" + evictPrimaryMetrics = "30327" + + evictBucket = "evictsrc" + evictMount = "evictmnt" +) + +var miniLogPaths []string + +func startMini(t *testing.T, dir string, args ...string) { + t.Helper() + require.NoError(t, os.MkdirAll(dir, 0755)) + logPath := filepath.Join(dir, "weed.log") + logFile, err := os.Create(logPath) + require.NoError(t, err) + miniLogPaths = append(miniLogPaths, logPath) + cmd := exec.Command(weedBinary, append([]string{"mini", + "-dir=" + dir, + "-s3.config=s3_config.json", + "-s3.allowDeleteBucketNotEmpty=true", + "-ip=127.0.0.1", "-ip.bind=127.0.0.1", + }, args...)...) + cmd.Stdout = logFile + cmd.Stderr = logFile + require.NoError(t, cmd.Start()) + t.Cleanup(func() { + cmd.Process.Kill() + cmd.Wait() + logFile.Close() + }) +} + +func waitForHTTP(t *testing.T, url string) { + t.Helper() + deadline := time.Now().Add(90 * time.Second) + for time.Now().Before(deadline) { + if resp, err := http.Get(url); err == nil { + resp.Body.Close() + return + } + time.Sleep(time.Second) + } + for _, p := range miniLogPaths { + if data, err := os.ReadFile(p); err == nil { + lines := strings.Split(string(data), "\n") + if len(lines) > 30 { + lines = lines[len(lines)-30:] + } + t.Logf("last lines of %s:\n%s", p, strings.Join(lines, "\n")) + } + } + t.Fatalf("timed out waiting for %s", url) +} + +func shellOn(t *testing.T, masterPort, command string) string { + t.Helper() + cmd := exec.Command(weedBinary, "shell", "-master=localhost:"+masterPort) + cmd.Stdin = strings.NewReader(command + "\nexit\n") + out, err := cmd.CombinedOutput() + require.NoErrorf(t, err, "shell %q failed: %s", command, out) + return stripLogs(string(out)) +} + +func chunkCountOn(t *testing.T, masterPort, path string) string { + meta := shellOn(t, masterPort, "fs.meta.cat "+path) + idx := strings.LastIndex(meta, "chunks ") + require.GreaterOrEqualf(t, idx, 0, "no chunk count in %s", meta) + return strings.Fields(meta[idx+len("chunks "):])[0] +} + +func volumeStats(t *testing.T, volumePort string) (size, garbage uint64) { + t.Helper() + resp, err := http.Get("http://localhost:" + volumePort + "/status") + require.NoError(t, err) + defer resp.Body.Close() + var status struct { + Volumes []struct { + Size uint64 `json:"Size"` + DeletedByteCount uint64 `json:"DeletedByteCount"` + } `json:"Volumes"` + } + require.NoError(t, json.NewDecoder(resp.Body).Decode(&status)) + for _, v := range status.Volumes { + size += v.Size + garbage += v.DeletedByteCount + } + return size, garbage +} + +func readViaFiler(t *testing.T, filerPort, path string) []byte { + t.Helper() + resp, err := http.Get("http://localhost:" + filerPort + path) + require.NoError(t, err) + defer resp.Body.Close() + require.Equal(t, http.StatusOK, resp.StatusCode, "read %s", path) + data, err := io.ReadAll(resp.Body) + require.NoError(t, err) + return data +} + +// TestRemoteCacheEvictUnderPressure fills a cache constrained to two small +// volumes with remote-mounted objects until writes fail, then verifies the +// filer evicts the oldest synced entry, vacuums the garbage, and a later read +// caches again. +func TestRemoteCacheEvictUnderPressure(t *testing.T) { + if testing.Short() { + t.Skip("spawns two weed mini clusters") + } + if _, err := os.Stat(weedBinary); err != nil { + t.Skipf("weed binary not found at %s; run make build-weed", weedBinary) + } + if isServerRunning(evictRemoteS3) || isServerRunning(evictPrimaryS3) { + t.Skip("eviction test ports are already in use") + } + + tmp := t.TempDir() + startMini(t, filepath.Join(tmp, "remote"), + "-s3.port=28334", "-master.port="+evictRemoteMaster, + "-filer.port="+evictRemoteFiler, "-volume.port="+evictRemoteVolume, + "-webdav.port="+evictRemoteWebdav, "-metricsPort="+evictRemoteMetrics) + waitForHTTP(t, evictRemoteS3) + startMini(t, filepath.Join(tmp, "primary"), + "-s3.port=28333", "-master.port="+evictPrimaryMaster, + "-filer.port="+evictPrimaryFiler, "-volume.port="+evictPrimaryVolume, + "-webdav.port="+evictPrimaryWebdav, "-metricsPort="+evictPrimaryMetrics, + "-volume.allowUntrustedRemoteEndpoints", "-filer.allowUntrustedRemoteEndpoints", + "-s3.allowUntrustedRemoteEndpoints", + "-master.volumeSizeLimitMB=32", "-volume.max=2", + "-filer.remoteCacheEvictThreshold=0.99") + waitForHTTP(t, evictPrimaryS3) + + remote := createS3Client(evictRemoteS3) + _, err := remote.CreateBucket(&s3.CreateBucketInput{Bucket: aws.String(evictBucket)}) + require.NoError(t, err) + for i := 0; i < 5; i++ { + data := make([]byte, 16*1024*1024) + for j := range data { + data[j] = byte(i + j%251) + } + _, err = remote.PutObject(&s3.PutObjectInput{ + Bucket: aws.String(evictBucket), + Key: aws.String(fmt.Sprintf("obj%d.bin", i)), + Body: bytes.NewReader(data), + }) + require.NoError(t, err) + } + + shellOn(t, evictPrimaryMaster, fmt.Sprintf( + "remote.configure -name=evictremote -type=s3 -s3.access_key=%s -s3.secret_key=%s -s3.endpoint=%s -s3.region=us-east-1", + accessKey, secretKey, evictRemoteS3)) + shellOn(t, evictPrimaryMaster, fmt.Sprintf( + "remote.mount -dir=/buckets/%s -remote=evictremote/%s -nonempty", evictMount, evictBucket)) + shellOn(t, evictPrimaryMaster, fmt.Sprintf("remote.meta.sync -dir=/buckets/%s", evictMount)) + time.Sleep(2 * time.Second) + + mount := "/buckets/" + evictMount + + // obj0 caches alone so it is the oldest evictable entry. + first := readViaFiler(t, evictPrimaryFiler, mount+"/obj0.bin") + require.Len(t, first, 16*1024*1024) + require.NotEqual(t, "0", chunkCountOn(t, evictPrimaryMaster, mount+"/obj0.bin"), "obj0 should be cached") + + // 4 x 16MB against ~64MB of capacity: the fills overrun, hit the + // capacity error path, and trigger eviction + vacuum. + var wg sync.WaitGroup + for i := 1; i <= 4; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + data := readViaFiler(t, evictPrimaryFiler, fmt.Sprintf("%s/obj%d.bin", mount, i)) + assert.Len(t, data, 16*1024*1024) + }(i) + } + wg.Wait() + + // Eviction must have dropped obj0's local chunks and vacuumed the + // tombstoned bytes, leaving only live data within the two volumes. + require.True(t, waitForCondition(t, func() bool { + size, garbage := volumeStats(t, evictPrimaryVolume) + return garbage == 0 && size <= 72*1024*1024 + }, 2*time.Minute, "volumes to be reclaimed by eviction+vacuum"), + "evicted chunks were not reclaimed") + + assert.Equal(t, "0", chunkCountOn(t, evictPrimaryMaster, mount+"/obj0.bin"), + "oldest cached entry should be evicted back to remote-only") + + // The cache self-heals: reading the evicted object re-caches it. Fills + // still running from the concurrent wave may evict it again, so retry + // the read until a commit sticks. + var again []byte + require.True(t, waitForCondition(t, func() bool { + again = readViaFiler(t, evictPrimaryFiler, mount+"/obj0.bin") + return chunkCountOn(t, evictPrimaryMaster, mount+"/obj0.bin") != "0" + }, 2*time.Minute, "obj0 to re-cache"), + "evicted object did not re-cache on read") + assert.Equal(t, first, again) +} diff --git a/weed/command/filer.go b/weed/command/filer.go index 97cd923aa..db17fa5fd 100644 --- a/weed/command/filer.go +++ b/weed/command/filer.go @@ -89,6 +89,7 @@ type FilerOptions struct { s3ConfigFile *string // optional path to static S3 identity config allowUntrustedRemoteEndpoints *bool + remoteCacheEvictThreshold *float64 // shutdownCtx, when non-nil, tells startFiler to gracefully shut down its // HTTP/gRPC servers once the ctx is cancelled. Used by integration tests // and by weed mini; nil for standalone weed filer. @@ -134,6 +135,7 @@ func init() { f.tusMaxSizeMB = cmdFiler.Flag.Int("tusMaxSizeMB", 5*1024, "maximum TUS upload size in MB") f.tusSessionExpiry = cmdFiler.Flag.Duration("tusSessionExpiry", 24*time.Hour, "incomplete TUS upload sessions are cleaned up after this duration, e.g. \"48h\", \"7h30m\"") f.allowUntrustedRemoteEndpoints = cmdFiler.Flag.Bool("allowUntrustedRemoteEndpoints", false, allowUntrustedRemoteEndpointsUsage) + f.remoteCacheEvictThreshold = cmdFiler.Flag.Float64("remoteCacheEvictThreshold", 0.9, "evict remote-cached objects (oldest first) when any volume disk exceeds this usage fraction; 0 disables") // start s3 on filer filerStartS3 = cmdFiler.Flag.Bool("s3", false, "whether to start S3 gateway") @@ -407,6 +409,7 @@ func (fo *FilerOptions) startFiler() { CredentialManager: credentialManager, AllowUntrustedRemoteEndpoints: *fo.allowUntrustedRemoteEndpoints, + RemoteCacheEvictThreshold: *fo.remoteCacheEvictThreshold, }) if nfs_err != nil { glog.Fatalf("Filer startup error: %v", nfs_err) diff --git a/weed/command/mini.go b/weed/command/mini.go index e5c2460dd..440aaa851 100644 --- a/weed/command/mini.go +++ b/weed/command/mini.go @@ -467,6 +467,7 @@ func initMiniFilerFlags() { miniFilerOptions.tusMaxSizeMB = cmdMini.Flag.Int("filer.tusMaxSizeMB", 5*1024, "maximum TUS upload size in MB") miniFilerOptions.tusSessionExpiry = cmdMini.Flag.Duration("filer.tusSessionExpiry", 24*time.Hour, "incomplete TUS upload sessions are cleaned up after this duration") miniFilerOptions.allowUntrustedRemoteEndpoints = cmdMini.Flag.Bool("filer.allowUntrustedRemoteEndpoints", false, allowUntrustedRemoteEndpointsUsage) + miniFilerOptions.remoteCacheEvictThreshold = cmdMini.Flag.Float64("filer.remoteCacheEvictThreshold", 0.9, "evict remote-cached objects (oldest first) when any volume disk exceeds this usage fraction; 0 disables") } // initMiniVolumeFlags initializes Volume server flag options diff --git a/weed/command/server.go b/weed/command/server.go index 6fbacae0d..da1385a08 100644 --- a/weed/command/server.go +++ b/weed/command/server.go @@ -136,6 +136,7 @@ func init() { filerOptions.tusMaxSizeMB = cmdServer.Flag.Int("filer.tusMaxSizeMB", 5*1024, "maximum TUS upload size in MB") filerOptions.tusSessionExpiry = cmdServer.Flag.Duration("filer.tusSessionExpiry", 24*time.Hour, "incomplete TUS upload sessions are cleaned up after this duration, e.g. \"48h\", \"7h30m\"") filerOptions.allowUntrustedRemoteEndpoints = cmdServer.Flag.Bool("filer.allowUntrustedRemoteEndpoints", false, allowUntrustedRemoteEndpointsUsage) + filerOptions.remoteCacheEvictThreshold = cmdServer.Flag.Float64("filer.remoteCacheEvictThreshold", 0.9, "evict remote-cached objects (oldest first) when any volume disk exceeds this usage fraction; 0 disables") serverOptions.v.port = cmdServer.Flag.Int("volume.port", 8080, "volume server http listen port") serverOptions.v.portGrpc = cmdServer.Flag.Int("volume.port.grpc", 0, "volume server grpc listen port") diff --git a/weed/filer/filer.go b/weed/filer/filer.go index d83486c86..e957c4d16 100644 --- a/weed/filer/filer.go +++ b/weed/filer/filer.go @@ -23,6 +23,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/pb/remote_pb" "google.golang.org/grpc" + "google.golang.org/protobuf/proto" "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" @@ -327,7 +328,7 @@ func (f *Filer) CreateEntry(ctx context.Context, entry *Entry, existing *Entry, return fmt.Errorf("%s: %w", entry.FullPath, filer_pb.ErrEntryAlreadyExists) } glog.V(4).InfofCtx(ctx, "UpdateEntry %s: old entry: %v", entry.FullPath, oldEntry.Name()) - if err := f.UpdateEntry(ctx, oldEntry, entry); err != nil { + if err := f.UpdateEntry(ctx, oldEntry, entry, isFromOtherCluster); err != nil { if errors.Is(err, filer_pb.ErrExistingIsDirectory) || errors.Is(err, filer_pb.ErrExistingIsFile) { glog.V(2).InfofCtx(ctx, "update entry %s: %v", entry.FullPath, err) } else { @@ -482,7 +483,7 @@ func (f *Filer) EnsureDirectoryEntry(ctx context.Context, dirPath util.FullPath, narrowed := existing.ShallowClone() narrowed.Mode = existing.Mode&^restorableModeBits | kept glog.V(1).InfofCtx(ctx, "restore directory %s: narrowing %v to %v", dirPath, existing.Mode, narrowed.Mode) - if err := f.UpdateEntry(ctx, existing, narrowed); err != nil { + if err := f.UpdateEntry(ctx, existing, narrowed, false); err != nil { return err } f.NotifyUpdateEvent(ctx, existing, narrowed, false, false, nil) @@ -515,7 +516,7 @@ func (f *Filer) EnsureDirectoryEntry(ctx context.Context, dirPath util.FullPath, return nil } -func (f *Filer) UpdateEntry(ctx context.Context, oldEntry, entry *Entry) (err error) { +func (f *Filer) UpdateEntry(ctx context.Context, oldEntry, entry *Entry, isFromOtherCluster bool) (err error) { if oldEntry != nil { entry.Attr.Crtime = oldEntry.Attr.Crtime if oldEntry.Attr.Inode != 0 { @@ -535,6 +536,17 @@ func (f *Filer) UpdateEntry(ctx context.Context, oldEntry, entry *Entry) (err er glog.V(2).InfofCtx(ctx, "existing %s is a file", oldEntry.FullPath) return fmt.Errorf("%s: %w", oldEntry.FullPath, filer_pb.ErrExistingIsFile) } + // A local write to a remote-backed entry leaves the copy on remote + // stale until a sync re-uploads it. Content changes that arrive + // without a fresh sync stamp are unsynced; clear the stamp so nothing + // mistakes the still-local chunks for a re-fetchable cache copy. + // Replicated updates carry the writer's authoritative stamp. + if !isFromOtherCluster && oldEntry.Remote != nil && entry.Remote != nil && + oldEntry.Remote.LastLocalSyncTsNs == entry.Remote.LastLocalSyncTsNs && + !chunksEqual(oldEntry.Chunks, entry.Chunks) { + entry.Remote = proto.Clone(entry.Remote).(*filer_pb.RemoteEntry) + entry.Remote.LastLocalSyncTsNs = 0 + } } if entry.Attr.Atime.IsZero() { entry.Attr.Atime = entryInitialAtime(entry.Attr) diff --git a/weed/filer/filer_deletion.go b/weed/filer/filer_deletion.go index f32b97a8b..5ca125a31 100644 --- a/weed/filer/filer_deletion.go +++ b/weed/filer/filer_deletion.go @@ -306,24 +306,29 @@ func (f *Filer) loopProcessingDeletion() { glog.V(0).Infof("deletion processor shutting down") return case <-ticker.C: - f.FileIdDeletionQueue.Consume(func(fileIds []string) { - for i := 0; i < len(fileIds); i += DeletionBatchSize { - end := i + DeletionBatchSize - if end > len(fileIds) { - end = len(fileIds) - } - toDeleteFileIds := fileIds[i:end] - f.processDeletionBatch(toDeleteFileIds, lookupFunc) - } - }) + f.FlushFileIdDeletionQueue(context.Background(), lookupFunc) } } } +func (f *Filer) FlushFileIdDeletionQueue(ctx context.Context, lookupFunc func([]string) (map[string]*operation.LookupResult, error)) (consumed []string) { + f.FileIdDeletionQueue.Consume(func(fileIds []string) { + consumed = fileIds + for i := 0; i < len(fileIds); i += DeletionBatchSize { + end := i + DeletionBatchSize + if end > len(fileIds) { + end = len(fileIds) + } + f.processDeletionBatch(ctx, fileIds[i:end], lookupFunc) + } + }) + return consumed +} + // processDeletionBatch handles deletion of a batch of file IDs and processes results. // It classifies errors into retryable and permanent categories, adds retryable failures // to the retry queue, and logs appropriate messages. -func (f *Filer) processDeletionBatch(toDeleteFileIds []string, lookupFunc func([]string) (map[string]*operation.LookupResult, error)) { +func (f *Filer) processDeletionBatch(ctx context.Context, toDeleteFileIds []string, lookupFunc func([]string) (map[string]*operation.LookupResult, error)) { // Deduplicate file IDs to prevent incorrect retry count increments for the same file ID within a single batch. uniqueFileIdsSlice := make([]string, 0, len(toDeleteFileIds)) processed := make(map[string]struct{}, len(toDeleteFileIds)) @@ -339,7 +344,7 @@ func (f *Filer) processDeletionBatch(toDeleteFileIds []string, lookupFunc func([ } // Delete files and classify outcomes - outcomes := deleteFilesAndClassify(f.GrpcDialOption, uniqueFileIdsSlice, lookupFunc) + outcomes := deleteFilesAndClassify(ctx, f.GrpcDialOption, uniqueFileIdsSlice, lookupFunc) // Process outcomes var successCount, notFoundCount, retryableErrorCount, permanentErrorCount int @@ -405,9 +410,9 @@ type deletionOutcome struct { } // deleteFilesAndClassify performs deletion and classifies outcomes for a list of file IDs -func deleteFilesAndClassify(grpcDialOption grpc.DialOption, fileIds []string, lookupFunc func([]string) (map[string]*operation.LookupResult, error)) map[string]deletionOutcome { +func deleteFilesAndClassify(ctx context.Context, grpcDialOption grpc.DialOption, fileIds []string, lookupFunc func([]string) (map[string]*operation.LookupResult, error)) map[string]deletionOutcome { // Perform deletion - results := operation.DeleteFileIdsWithLookupVolumeId(grpcDialOption, fileIds, lookupFunc) + results := operation.DeleteFileIdsWithLookupVolumeId(ctx, grpcDialOption, fileIds, lookupFunc) // Group results by file ID to handle multiple results for replicated volumes resultsByFileId := make(map[string][]*volume_server_pb.DeleteResult) @@ -545,7 +550,7 @@ func (f *Filer) processRetryBatch(readyItems []*DeletionRetryItem, lookupFunc fu } // Delete files and classify outcomes - outcomes := deleteFilesAndClassify(f.GrpcDialOption, fileIds, lookupFunc) + outcomes := deleteFilesAndClassify(context.Background(), f.GrpcDialOption, fileIds, lookupFunc) // Process outcomes - iterate over readyItems to ensure all items are accounted for var successCount, notFoundCount, retryCount, permanentErrorCount int diff --git a/weed/filer/filer_inode_test.go b/weed/filer/filer_inode_test.go index 42990a6fb..575332ca0 100644 --- a/weed/filer/filer_inode_test.go +++ b/weed/filer/filer_inode_test.go @@ -173,7 +173,7 @@ func TestUpdateEntryPreservesExistingInode(t *testing.T) { }, } - err := f.UpdateEntry(context.Background(), original, updated) + err := f.UpdateEntry(context.Background(), original, updated, false) require.Error(t, err) updated = &Entry{ @@ -182,7 +182,7 @@ func TestUpdateEntryPreservesExistingInode(t *testing.T) { Mode: 0o600, }, } - err = f.UpdateEntry(context.Background(), original, updated) + err = f.UpdateEntry(context.Background(), original, updated, false) require.NoError(t, err) stored, findErr := store.FindEntry(context.Background(), original.FullPath) @@ -208,7 +208,7 @@ func TestUpdateEntryBackfillsMissingLegacyInode(t *testing.T) { Mode: 0o640, }, } - err := f.UpdateEntry(context.Background(), original, updated) + err := f.UpdateEntry(context.Background(), original, updated, false) require.NoError(t, err) stored, findErr := store.FindEntry(context.Background(), original.FullPath) diff --git a/weed/filer/filer_remote_evict.go b/weed/filer/filer_remote_evict.go new file mode 100644 index 000000000..e29548bfb --- /dev/null +++ b/weed/filer/filer_remote_evict.go @@ -0,0 +1,89 @@ +package filer + +import ( + "context" + "sort" + "time" + + "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +func chunksEqual(a, b []*filer_pb.FileChunk) bool { + if len(a) != len(b) { + return false + } + for i := range a { + if a[i].GetFileIdString() != b[i].GetFileIdString() || + a[i].Size != b[i].Size || + a[i].Offset != b[i].Offset || + a[i].ModifiedTsNs != b[i].ModifiedTsNs { + return false + } + } + return true +} + +// IsEvictableRemoteEntry reports whether an entry's local chunks are backed by +// a synchronized remote copy, mirroring the checks remote.uncache applies: +// remote-backed, chunks present, and not ahead of the remote version. +func IsEvictableRemoteEntry(entry *Entry) bool { + if entry.IsDirectory() || entry.Remote == nil { + return false + } + if entry.Remote.LastLocalSyncTsNs <= 0 || len(entry.GetChunks()) == 0 { + return false + } + if entry.Remote.LastLocalSyncTsNs < entry.Mtime.UnixNano() { + return false + } + return true +} + +// ListEvictableRemoteEntries walks every mounted directory and returns +// synchronized remote entries holding local chunks, oldest cached first. +func (f *Filer) ListEvictableRemoteEntries(ctx context.Context, mounts []util.FullPath, minCacheAge time.Duration) (out []*Entry) { + cutoffNs := time.Now().UnixNano() - minCacheAge.Nanoseconds() + for _, dir := range mounts { + if err := f.collectEvictableRemoteEntries(ctx, dir, cutoffNs, &out); err != nil { + glog.WarningfCtx(ctx, "list evictable remote entries under %s: %v", dir, err) + } + } + sort.Slice(out, func(i, j int) bool { + return out[i].Remote.LastLocalSyncTsNs < out[j].Remote.LastLocalSyncTsNs + }) + return out +} + +func (f *Filer) collectEvictableRemoteEntries(ctx context.Context, dir util.FullPath, cutoffNs int64, out *[]*Entry) error { + startFileName := "" + for { + var subDirs []util.FullPath + lastFileName, err := f.Store.ListDirectoryEntries(ctx, dir, startFileName, false, 1024, func(entry *Entry) (bool, error) { + if err := ctx.Err(); err != nil { + return false, err + } + if entry.IsDirectory() { + subDirs = append(subDirs, entry.FullPath) + return true, nil + } + if entry.Remote != nil && entry.Remote.LastLocalSyncTsNs > 0 && entry.Remote.LastLocalSyncTsNs <= cutoffNs && IsEvictableRemoteEntry(entry) { + *out = append(*out, entry) + } + return true, nil + }) + if err != nil { + return err + } + for _, subDir := range subDirs { + if err := f.collectEvictableRemoteEntries(ctx, subDir, cutoffNs, out); err != nil { + return err + } + } + if lastFileName == "" { + return nil + } + startFileName = lastFileName + } +} diff --git a/weed/filer/filer_remote_evict_test.go b/weed/filer/filer_remote_evict_test.go new file mode 100644 index 000000000..a1c8717cc --- /dev/null +++ b/weed/filer/filer_remote_evict_test.go @@ -0,0 +1,164 @@ +package filer + +import ( + "context" + "os" + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/util" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func remoteCachedEntry(path string, syncTsNs int64, mtime time.Time, chunks int) *Entry { + entry := &Entry{ + FullPath: util.FullPath(path), + Attr: Attr{ + Mtime: mtime, + Crtime: mtime, + Mode: 0644, + Uid: 1, + Gid: 1, + Mime: "application/octet-stream", + Md5: nil, + FileSize: 0, + }, + Remote: &filer_pb.RemoteEntry{ + RemoteMtime: mtime.Unix(), + LastLocalSyncTsNs: syncTsNs, + RemoteETag: "etag", + RemoteSize: 100, + }, + } + for i := 0; i < chunks; i++ { + entry.Chunks = append(entry.Chunks, &filer_pb.FileChunk{FileId: "1,01637037d6", Size: 100}) + } + return entry +} + +func TestIsEvictableRemoteEntry(t *testing.T) { + now := time.Now() + synced := now.Add(-time.Hour).UnixNano() + mtime := now.Add(-time.Hour) + + tests := []struct { + name string + entry *Entry + expected bool + }{ + {"directory", &Entry{FullPath: "/d", Attr: Attr{Mode: os.ModeDir | 0755}}, false}, + {"local only, no remote", &Entry{FullPath: "/f", Attr: Attr{Mtime: mtime, Mode: 0644}}, false}, + {"remote only, never cached", remoteCachedEntry("/f", 0, mtime, 0), false}, + {"remote, chunks cleared", remoteCachedEntry("/f", synced, mtime, 0), false}, + {"dirty, newer than remote", remoteCachedEntry("/f", synced, now, 1), false}, + {"same-second write after sync", remoteCachedEntry("/f", + now.Truncate(time.Second).Add(100*time.Millisecond).UnixNano(), + now.Truncate(time.Second).Add(900*time.Millisecond), 1), false}, + {"same-second sync after write", remoteCachedEntry("/f", + now.Truncate(time.Second).Add(900*time.Millisecond).UnixNano(), + now.Truncate(time.Second).Add(100*time.Millisecond), 1), true}, + {"synced and cached", remoteCachedEntry("/f", synced, mtime, 1), true}, + } + for _, tt := range tests { + assert.Equal(t, tt.expected, IsEvictableRemoteEntry(tt.entry), tt.name) + } +} + +func TestListEvictableRemoteEntries(t *testing.T) { + now := time.Now() + mtime := now.Add(-2 * time.Hour) + oldSync := now.Add(-time.Hour).UnixNano() + newSync := now.Add(-time.Second).UnixNano() + + store := newStubFilerStore() + f := newTestFiler(t, store, NewFilerRemoteStorage()) + mounts := []util.FullPath{"/buckets/mybucket"} + + seed := func(e *Entry) { + require.NoError(t, f.CreateEntry(context.Background(), e, nil, false, false, nil, false, 255)) + } + seed(remoteCachedEntry("/buckets/mybucket/old.bin", oldSync, mtime, 1)) + seed(remoteCachedEntry("/buckets/mybucket/fresh.bin", newSync, mtime, 1)) + seed(remoteCachedEntry("/buckets/mybucket/remoteonly.bin", 0, mtime, 0)) + seed(remoteCachedEntry("/buckets/mybucket/dirty.bin", oldSync, now, 1)) + seed(&Entry{FullPath: "/buckets/mybucket/plain.bin", Attr: Attr{Mtime: mtime, Mode: 0644}}) + + store.entries["/buckets/mybucket/sub"] = &Entry{ + FullPath: "/buckets/mybucket/sub", + Attr: Attr{Mode: os.ModeDir | 0755, Mtime: mtime}, + } + seed(remoteCachedEntry("/buckets/mybucket/sub/nested.bin", oldSync-1000, mtime, 1)) + + got := f.ListEvictableRemoteEntries(context.Background(), mounts, 0) + var paths []string + for _, e := range got { + paths = append(paths, string(e.FullPath)) + } + assert.Equal(t, []string{ + "/buckets/mybucket/sub/nested.bin", + "/buckets/mybucket/old.bin", + "/buckets/mybucket/fresh.bin", + }, paths, "oldest LastLocalSyncTsNs first; remote-only, dirty, local entries skipped") + + got = f.ListEvictableRemoteEntries(context.Background(), mounts, 30*time.Second) + paths = paths[:0] + for _, e := range got { + paths = append(paths, string(e.FullPath)) + } + assert.Equal(t, []string{ + "/buckets/mybucket/sub/nested.bin", + "/buckets/mybucket/old.bin", + }, paths, "minCacheAge excludes freshly cached entries") +} + +func TestUpdateEntryInvalidatesStaleSyncStamp(t *testing.T) { + now := time.Now() + synced := now.Add(-time.Hour).UnixNano() + mtime := now.Add(-time.Hour) + path := "/buckets/mybucket/a.bin" + + newFiler := func() *Filer { + f := newTestFiler(t, newStubFilerStore(), NewFilerRemoteStorage()) + require.NoError(t, f.CreateEntry(context.Background(), remoteCachedEntry(path, synced, mtime, 1), nil, false, false, nil, false, 255)) + return f + } + syncStamp := func(f *Filer) int64 { + entry, err := f.FindEntry(context.Background(), util.FullPath(path)) + require.NoError(t, err) + return entry.Remote.LastLocalSyncTsNs + } + + t.Run("local chunk change clears stamp", func(t *testing.T) { + f := newFiler() + update := remoteCachedEntry(path, synced, mtime, 1) + update.Chunks[0].FileId = "2,01637037d6" + require.NoError(t, f.CreateEntry(context.Background(), update, nil, false, false, nil, false, 255)) + assert.Zero(t, syncStamp(f)) + }) + + t.Run("metadata-only update keeps stamp", func(t *testing.T) { + f := newFiler() + update := remoteCachedEntry(path, synced, mtime, 1) + update.Attr.Mime = "text/plain" + require.NoError(t, f.CreateEntry(context.Background(), update, nil, false, false, nil, false, 255)) + assert.Equal(t, synced, syncStamp(f)) + }) + + t.Run("fresh sync stamp survives chunk change", func(t *testing.T) { + f := newFiler() + update := remoteCachedEntry(path, now.UnixNano(), mtime, 1) + update.Chunks[0].FileId = "2,01637037d6" + require.NoError(t, f.CreateEntry(context.Background(), update, nil, false, false, nil, false, 255)) + assert.Equal(t, now.UnixNano(), syncStamp(f)) + }) + + t.Run("replicated update keeps authoritative stamp", func(t *testing.T) { + f := newFiler() + update := remoteCachedEntry(path, synced, mtime, 1) + update.Chunks[0].FileId = "2,01637037d6" + require.NoError(t, f.CreateEntry(context.Background(), update, nil, false, true, nil, false, 255)) + assert.Equal(t, synced, syncStamp(f)) + }) +} diff --git a/weed/operation/delete_content.go b/weed/operation/delete_content.go index 5028fbf48..297f66515 100644 --- a/weed/operation/delete_content.go +++ b/weed/operation/delete_content.go @@ -45,11 +45,11 @@ func DeleteFileIds(masterFn GetMasterFn, usePublicUrl bool, grpcDialOption grpc. return } - return DeleteFileIdsWithLookupVolumeId(grpcDialOption, fileIds, lookupFunc) + return DeleteFileIdsWithLookupVolumeId(context.Background(), grpcDialOption, fileIds, lookupFunc) } -func DeleteFileIdsWithLookupVolumeId(grpcDialOption grpc.DialOption, fileIds []string, lookupFunc func(vid []string) (map[string]*LookupResult, error)) []*volume_server_pb.DeleteResult { +func DeleteFileIdsWithLookupVolumeId(ctx context.Context, grpcDialOption grpc.DialOption, fileIds []string, lookupFunc func(vid []string) (map[string]*LookupResult, error)) []*volume_server_pb.DeleteResult { var ret []*volume_server_pb.DeleteResult @@ -117,7 +117,7 @@ func DeleteFileIdsWithLookupVolumeId(grpcDialOption grpc.DialOption, fileIds []s go func(server pb.ServerAddress, fidList []string) { defer wg.Done() - resultChan <- DeleteFileIdsAtOneVolumeServer(server, grpcDialOption, fidList, false) + resultChan <- DeleteFileIdsAtOneVolumeServer(ctx, server, grpcDialOption, fidList, false) }(server, fidList) } @@ -133,7 +133,7 @@ func DeleteFileIdsWithLookupVolumeId(grpcDialOption grpc.DialOption, fileIds []s // DeleteFileIdsAtOneVolumeServer deletes a list of files that is on one volume server via gRpc // Returns individual results for each file ID. Check result.Error for per-file failures. -func DeleteFileIdsAtOneVolumeServer(volumeServer pb.ServerAddress, grpcDialOption grpc.DialOption, fileIds []string, includeCookie bool) []*volume_server_pb.DeleteResult { +func DeleteFileIdsAtOneVolumeServer(ctx context.Context, volumeServer pb.ServerAddress, grpcDialOption grpc.DialOption, fileIds []string, includeCookie bool) []*volume_server_pb.DeleteResult { var ret []*volume_server_pb.DeleteResult @@ -144,7 +144,7 @@ func DeleteFileIdsAtOneVolumeServer(volumeServer pb.ServerAddress, grpcDialOptio SkipCookieCheck: !includeCookie, } - resp, err := volumeServerClient.BatchDelete(context.Background(), req) + resp, err := volumeServerClient.BatchDelete(ctx, req) // fmt.Printf("deleted %v %v: %v\n", fileIds, err, resp) diff --git a/weed/s3api/s3api_object_handlers_put.go b/weed/s3api/s3api_object_handlers_put.go index 2a8f9259f..53f7fe027 100644 --- a/weed/s3api/s3api_object_handlers_put.go +++ b/weed/s3api/s3api_object_handlers_put.go @@ -2738,7 +2738,7 @@ func (s3a *S3ApiServer) deleteOrphanedChunks(chunks []*filer_pb.FileChunk) { } // Attempt deletion using the operation package's batch delete with custom lookup - deleteResults := operation.DeleteFileIdsWithLookupVolumeId(s3a.option.GrpcDialOption, fileIds, lookupFunc) + deleteResults := operation.DeleteFileIdsWithLookupVolumeId(context.Background(), s3a.option.GrpcDialOption, fileIds, lookupFunc) // Log results - track successes and failures successCount := 0 diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index 9887f6bce..85c73b02e 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -518,7 +518,7 @@ func (fs *FilerServer) applyObjectMutation(ctx context.Context, m *filer_pb.Obje if m.TouchMtime { newEntry.Attr.Mtime = time.Now() } - if err := fs.filer.UpdateEntry(ctx, oldEntry, newEntry); err != nil { + if err := fs.filer.UpdateEntry(ctx, oldEntry, newEntry, fromOtherCluster); err != nil { return err } // Emit the metadata event so the update replicates and subscribers see it, @@ -623,7 +623,7 @@ func (fs *FilerServer) applyRecomputeLatest(ctx context.Context, m *filer_pb.Obj } } - if err := fs.filer.UpdateEntry(ctx, oldPointer, pointer); err != nil { + if err := fs.filer.UpdateEntry(ctx, oldPointer, pointer, fromOtherCluster); err != nil { return err } // Replicate the recomputed pointer to peer filers and subscribers. Without @@ -653,7 +653,7 @@ func (fs *FilerServer) applyRecomputeLatest(ctx context.Context, m *filer_pb.Obj priorEntry.Extended = make(map[string][]byte) } priorEntry.Extended[rc.DemoteKey] = rc.DemoteValue - if err := fs.filer.UpdateEntry(ctx, oldPrior, priorEntry); err != nil { + if err := fs.filer.UpdateEntry(ctx, oldPrior, priorEntry, fromOtherCluster); err != nil { return err } fs.filer.NotifyUpdateEvent(ctx, oldPrior, priorEntry, false, fromOtherCluster, signatures) @@ -746,7 +746,7 @@ func (fs *FilerServer) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntr ctx, eventSink := filer.WithMetadataEventSink(ctx) resp := &filer_pb.UpdateEntryResponse{LogTsNs: logTsNs, LogSignature: fs.filer.Signature} - if err = fs.filer.UpdateEntry(ctx, entry, newEntry); err == nil { + if err = fs.filer.UpdateEntry(ctx, entry, newEntry, req.IsFromOtherCluster); err == nil { fs.filer.DeleteChunksNotRecursive(garbage) fs.filer.NotifyUpdateEvent(ctx, entry, newEntry, true, req.IsFromOtherCluster, req.Signatures) diff --git a/weed/server/filer_grpc_server_remote.go b/weed/server/filer_grpc_server_remote.go index 4c0589ad9..90f4b93b9 100644 --- a/weed/server/filer_grpc_server_remote.go +++ b/weed/server/filer_grpc_server_remote.go @@ -269,6 +269,14 @@ func (fs *FilerServer) doCacheRemoteObjectToLocalCluster(ctx context.Context, re 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 } diff --git a/weed/server/filer_server.go b/weed/server/filer_server.go index b1c859d11..2b2eec443 100644 --- a/weed/server/filer_server.go +++ b/weed/server/filer_server.go @@ -92,6 +92,9 @@ type FilerOption struct { // AllowUntrustedRemoteEndpoints lets a read of a remote-only entry dial a // mounted endpoint that resolves to a loopback / private / metadata host. AllowUntrustedRemoteEndpoints bool + // RemoteCacheEvictThreshold is the disk usage fraction at which the filer + // evicts remote-mounted cached chunks; 0 disables eviction. + RemoteCacheEvictThreshold float64 } type FilerServer struct { @@ -121,6 +124,15 @@ type FilerServer struct { // deduplicates concurrent remote object caching operations remoteCacheGroup singleflight.Group + // serializes remote-cache eviction passes; lastVacuum rate-limits the + // compaction trigger that reclaims evicted chunks. + remoteCacheEvictMu sync.Mutex + remoteCacheLastVacuum atomic.Pointer[time.Time] + remoteCacheEvictCtx context.Context + remoteCacheEvictCancel context.CancelFunc + remoteCachePendingVidsMu sync.Mutex + remoteCachePendingVids map[uint32]int + recentCopyRequestsMu sync.Mutex recentCopyRequests map[string]recentCopyRequest @@ -209,6 +221,7 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption) fs.startPosixLockSweeper() fs.mountPeerRegistry = filer.NewMountPeerRegistry() go fs.runMountPeerRegistrySweeper() + fs.remoteCacheEvictCtx, fs.remoteCacheEvictCancel = context.WithCancel(context.Background()) option.Masters.RefreshBySrvIfAvailable() if len(option.Masters.GetInstances()) == 0 { @@ -239,6 +252,7 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption) fs.filer.RemoteStorage.SetConfValidator(func(ctx context.Context, conf *remote_pb.RemoteConf) error { return ValidateRemoteConfForLoad(ctx, conf, option.AllowUntrustedRemoteEndpoints) }) + go fs.runRemoteCacheEviction() // we do not support IP whitelist right now https://github.com/seaweedfs/seaweedfs/issues/7094 if v.GetString("guard.white_list") != "" { glog.Warningf("filer: guard.white_list is configured but the IP whitelist feature is currently disabled. See https://github.com/seaweedfs/seaweedfs/issues/7094") @@ -363,6 +377,9 @@ func (fs *FilerServer) Shutdown() { if fs.posixLockSweeperStop != nil { close(fs.posixLockSweeperStop) } + if fs.remoteCacheEvictCancel != nil { + fs.remoteCacheEvictCancel() + } fs.filer.Shutdown() } diff --git a/weed/server/filer_server_remote_evict.go b/weed/server/filer_server_remote_evict.go new file mode 100644 index 000000000..1c11f5d46 --- /dev/null +++ b/weed/server/filer_server_remote_evict.go @@ -0,0 +1,300 @@ +package weed_server + +import ( + "context" + "strings" + "time" + + "github.com/seaweedfs/seaweedfs/weed/filer" + "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/pb/master_pb" + "github.com/seaweedfs/seaweedfs/weed/stats" + "github.com/seaweedfs/seaweedfs/weed/storage/needle" + "github.com/seaweedfs/seaweedfs/weed/util" + "google.golang.org/protobuf/proto" +) + +const ( + remoteCacheEvictInterval = 30 * time.Second + remoteCacheEvictMinAge = time.Minute + remoteCacheVacuumCooldown = time.Minute + remoteCacheMasterRpcTime = 30 * time.Second + remoteCachePendingVidAttempts = 10 +) + +// uncacheRemoteEntry drops the local chunks of one remote-mounted entry, the +// same state transition remote.uncache applies through UpdateEntry. Cleared +// chunks go to the deletion queue and are reclaimed by the next compaction. +// When vids is set, only entries holding chunks on those volumes count toward +// the freed bytes, and entries contributing nothing are left untouched. +func (fs *FilerServer) uncacheRemoteEntry(ctx context.Context, fullPath util.FullPath, minCacheAge time.Duration, vids map[uint32]struct{}) (freedBytes int64, err error) { + pathLock := fs.entryLockTable.AcquireLock("uncacheRemoteEntry", fullPath, util.ExclusiveLock) + defer fs.entryLockTable.ReleaseLock(fullPath, pathLock) + + current, err := fs.filer.FindEntry(ctx, fullPath) + if err != nil { + return 0, err + } + if !filer.IsEvictableRemoteEntry(current) { + return 0, nil + } + if time.Since(time.Unix(0, current.Remote.LastLocalSyncTsNs)) < minCacheAge { + return 0, nil + } + + freedBytes = remoteEntryBytesOnVids(current, vids) + if freedBytes == 0 { + return 0, nil + } + + newEntry := current.ShallowClone() + newEntry.Chunks = nil + newEntry.Remote = proto.Clone(current.Remote).(*filer_pb.RemoteEntry) + newEntry.Remote.LastLocalSyncTsNs = 0 + + if err := fs.filer.CreateEntry(ctx, newEntry, current, false, false, nil, true, fs.filer.MaxFilenameLength); err != nil { + return 0, err + } + fileIds := make([]string, 0, len(current.Chunks)) + for _, chunk := range current.Chunks { + fileIds = append(fileIds, chunk.GetFileIdString()) + } + fs.notePendingRemoteCacheVids(fileIds) + stats.RemoteCacheEvictedCounter.Inc() + glog.V(1).InfofCtx(ctx, "uncacheRemoteEntry %s freed %d bytes", fullPath, freedBytes) + return freedBytes, nil +} + +// remoteEntryBytesOnVids sums the entry's chunk bytes on the given volumes; a +// nil set counts the whole object. +func remoteEntryBytesOnVids(entry *filer.Entry, vids map[uint32]struct{}) int64 { + if vids == nil { + return int64(entry.Size()) + } + var bytes int64 + for _, chunk := range entry.Chunks { + fid, err := needle.ParseFileIdFromString(chunk.GetFileIdString()) + if err != nil { + continue + } + if _, ok := vids[uint32(fid.VolumeId)]; ok { + bytes += int64(chunk.Size) + } + } + return bytes +} + +// evictRemoteCachedEntries drops local chunks of remote-mounted entries +// oldest-cached first until bytesNeeded is met or candidates run out. The +// first pass honors a minimum cache age so a just-fetched hot object is not +// dropped under a reader; when aged candidates cannot cover the request a +// second pass accepts any synchronized cached entry. When pressuredVids is +// set, only bytes on those volumes count and entries elsewhere are skipped. +func (fs *FilerServer) evictRemoteCachedEntries(ctx context.Context, bytesNeeded int64, pressuredVids map[uint32]struct{}) (freed int64) { + if fs.filer.RemoteStorage == nil { + return 0 + } + mounts := fs.filer.RemoteStorage.MountedDirectories() + for _, minCacheAge := range []time.Duration{remoteCacheEvictMinAge, 0} { + for _, entry := range fs.filer.ListEvictableRemoteEntries(ctx, mounts, minCacheAge) { + if bytesNeeded > 0 && freed >= bytesNeeded { + return freed + } + n, err := fs.uncacheRemoteEntry(ctx, entry.FullPath, minCacheAge, pressuredVids) + if err != nil { + glog.WarningfCtx(ctx, "evict remote cache %s: %v", entry.FullPath, err) + continue + } + freed += n + } + } + return freed +} + +// remoteCacheDiskPressure reports per-disk usage across the cluster: the total +// bytes to reclaim and the volumes hosted on disks over the eviction threshold. +func (fs *FilerServer) remoteCacheDiskPressure(ctx context.Context) (bytesToFree int64, pressuredVids map[uint32]struct{}, over bool) { + threshold := fs.option.RemoteCacheEvictThreshold + if threshold <= 0 { + return 0, nil, false + } + pressuredVids = make(map[uint32]struct{}) + rpcCtx, cancel := context.WithTimeout(ctx, remoteCacheMasterRpcTime) + defer cancel() + err := fs.filer.MasterClient.WithClient(rpcCtx, false, func(client master_pb.SeaweedClient) error { + resp, err := client.VolumeList(rpcCtx, &master_pb.VolumeListRequest{}) + if err != nil { + return err + } + for _, dc := range resp.TopologyInfo.DataCenterInfos { + for _, rack := range dc.RackInfos { + for _, dn := range rack.DataNodeInfos { + for _, disk := range dn.DiskInfos { + for _, pd := range disk.SplitByPhysicalDisk() { + if pd.DiskTotalBytes == 0 { + continue + } + used := pd.DiskTotalBytes - pd.DiskFreeBytes + if float64(used) < float64(pd.DiskTotalBytes)*threshold { + continue + } + over = true + bytesToFree += int64(used) - int64(float64(pd.DiskTotalBytes)*threshold*0.95) + for _, vi := range pd.VolumeInfos { + pressuredVids[vi.Id] = struct{}{} + } + } + } + } + } + } + return nil + }) + if err != nil { + glog.WarningfCtx(ctx, "remote cache disk pressure check: %v", err) + return 0, nil, false + } + return bytesToFree, pressuredVids, over +} + +// flushAndVacuumRemoteCacheVolumes forces the deletion queue down to the volume +// servers so fresh tombstones land, then compacts the volumes carrying them. +// The flush survives shutdown cancellation: an interrupted delete would be +// requeued into the in-memory retry queue that exits with the process, +// stranding bytes whose metadata the eviction already dropped. +func (fs *FilerServer) flushAndVacuumRemoteCacheVolumes(ctx context.Context) { + if len(fs.pendingRemoteCacheVids()) == 0 { + return + } + flushCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), remoteCacheMasterRpcTime) + fs.filer.FlushFileIdDeletionQueue(flushCtx, filer.LookupByMasterClientFn(fs.filer.MasterClient)) + cancel() + fs.vacuumPendingRemoteCacheVids(ctx) +} + +// vacuumPendingRemoteCacheVids compacts every volume still owed a vacuum. The +// master does not report whether a request actually compacted, so a vid keeps +// roughly ten minutes of attempts: tombstones that land late or compaction +// losses to another vacuum are retried by later passes instead of stranded. +func (fs *FilerServer) vacuumPendingRemoteCacheVids(ctx context.Context) { + pending := fs.pendingRemoteCacheVids() + if len(pending) == 0 { + return + } + if last := fs.remoteCacheLastVacuum.Load(); last != nil && time.Since(*last) < remoteCacheVacuumCooldown { + return + } + now := time.Now() + fs.remoteCacheLastVacuum.Store(&now) + if err := fs.filer.MasterClient.WithClient(ctx, false, func(client master_pb.SeaweedClient) error { + for vid := range pending { + vCtx, cancel := context.WithTimeout(ctx, remoteCacheMasterRpcTime) + _, err := client.VacuumVolume(vCtx, &master_pb.VacuumVolumeRequest{VolumeId: vid}) + cancel() + if err != nil { + glog.WarningfCtx(ctx, "remote cache vacuum volume %d: %v", vid, err) + continue + } + fs.completePendingRemoteCacheVid(vid) + } + return nil + }); err != nil { + glog.WarningfCtx(ctx, "remote cache vacuum: %v", err) + } +} + +func (fs *FilerServer) notePendingRemoteCacheVids(fileIds []string) { + fs.remoteCachePendingVidsMu.Lock() + defer fs.remoteCachePendingVidsMu.Unlock() + if fs.remoteCachePendingVids == nil { + fs.remoteCachePendingVids = make(map[uint32]int) + } + for _, fid := range fileIds { + if parsed, err := needle.ParseFileIdFromString(fid); err == nil { + fs.remoteCachePendingVids[uint32(parsed.VolumeId)] = remoteCachePendingVidAttempts + } + } +} + +func (fs *FilerServer) pendingRemoteCacheVids() map[uint32]struct{} { + fs.remoteCachePendingVidsMu.Lock() + defer fs.remoteCachePendingVidsMu.Unlock() + out := make(map[uint32]struct{}, len(fs.remoteCachePendingVids)) + for vid := range fs.remoteCachePendingVids { + out[vid] = struct{}{} + } + return out +} + +func (fs *FilerServer) completePendingRemoteCacheVid(vid uint32) { + fs.remoteCachePendingVidsMu.Lock() + defer fs.remoteCachePendingVidsMu.Unlock() + if fs.remoteCachePendingVids[vid] <= 1 { + delete(fs.remoteCachePendingVids, vid) + } else { + fs.remoteCachePendingVids[vid]-- + } +} + +// reclaimRemoteCacheSpace evicts remote-cached content and compacts volumes to +// release disk space under capacity pressure. A pass already in flight is +// enough; callers that would queue behind it just fall back to remote reads. +func (fs *FilerServer) reclaimRemoteCacheSpace(ctx context.Context, bytesNeeded int64, pressuredVids map[uint32]struct{}) { + if !fs.remoteCacheEvictMu.TryLock() { + return + } + defer fs.remoteCacheEvictMu.Unlock() + freed := fs.evictRemoteCachedEntries(ctx, bytesNeeded, pressuredVids) + if freed > 0 { + glog.V(0).InfofCtx(ctx, "remote cache eviction freed %d bytes", freed) + } + fs.flushAndVacuumRemoteCacheVolumes(ctx) +} + +func isRemoteCacheCapacityError(err error) bool { + msg := err.Error() + return strings.Contains(msg, "writable volumes") || + strings.Contains(msg, "free volumes") || + strings.Contains(msg, "no space left") || + strings.Contains(msg, "out of space") +} + +// runRemoteCacheEviction periodically evicts remote-cached entries once any +// disk crosses the configured usage threshold, with a vacuum pass to reclaim +// the deleted chunks. +func (fs *FilerServer) evictCtx() context.Context { + if fs.remoteCacheEvictCtx == nil { + return context.Background() + } + return fs.remoteCacheEvictCtx +} + +func (fs *FilerServer) runRemoteCacheEviction() { + if fs.option.RemoteCacheEvictThreshold <= 0 || fs.remoteCacheEvictCtx == nil { + return + } + ctx := fs.remoteCacheEvictCtx + ticker := time.NewTicker(remoteCacheEvictInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + if fs.remoteCacheEvictMu.TryLock() { + fs.vacuumPendingRemoteCacheVids(ctx) + fs.remoteCacheEvictMu.Unlock() + } + if fs.filer.RemoteStorage == nil || len(fs.filer.RemoteStorage.MountedDirectories()) == 0 { + continue + } + bytesToFree, pressuredVids, over := fs.remoteCacheDiskPressure(ctx) + if !over { + continue + } + glog.V(0).Infof("remote cache: disk usage over %.0f%%, evicting %d bytes", fs.option.RemoteCacheEvictThreshold*100, bytesToFree) + fs.reclaimRemoteCacheSpace(ctx, bytesToFree, pressuredVids) + } +} diff --git a/weed/shell/command_fs_merge_volumes.go b/weed/shell/command_fs_merge_volumes.go index b992f811a..59aceecc2 100644 --- a/weed/shell/command_fs_merge_volumes.go +++ b/weed/shell/command_fs_merge_volumes.go @@ -339,7 +339,7 @@ func deleteOrphanedNeedles(commandEnv *CommandEnv, entryPath util.FullPath, need continue } for _, loc := range locations { - results := operation.DeleteFileIdsAtOneVolumeServer(loc.ServerAddress(), commandEnv.option.GrpcDialOption, fids, includeCookie) + results := operation.DeleteFileIdsAtOneVolumeServer(context.Background(), loc.ServerAddress(), commandEnv.option.GrpcDialOption, fids, includeCookie) // Summarize per server: an unreachable volume server returns one // error per needle, which for manifest-heavy files can mean // hundreds of near-identical lines. Keep the first error as the diff --git a/weed/shell/command_remote_uncache.go b/weed/shell/command_remote_uncache.go index d3d4371b3..e7281c791 100644 --- a/weed/shell/command_remote_uncache.go +++ b/weed/shell/command_remote_uncache.go @@ -102,7 +102,7 @@ func (c *commandRemoteUncache) uncacheContentData(commandEnv *CommandEnv, writer return true } - if entry.RemoteEntry.LastLocalSyncTsNs/1e9 < entry.Attributes.Mtime { + if entry.RemoteEntry.LastLocalSyncTsNs < entry.Attributes.Mtime*1e9+int64(entry.Attributes.MtimeNs) { return true // should not uncache an entry that is not synchronized with remote } diff --git a/weed/shell/command_volume_check_disk.go b/weed/shell/command_volume_check_disk.go index 79189eb30..d81929065 100644 --- a/weed/shell/command_volume_check_disk.go +++ b/weed/shell/command_volume_check_disk.go @@ -759,6 +759,7 @@ func (vcd *volumeCheckDisk) doVolumeCheckDisk(minuend, subtrahend *needle_map.Me vcd.writeVerbose("delete %s %s => %s", needleValue.Key.FileId(source.info.Id), source.location.dataNode.Id, target.location.dataNode.Id) } deleteResults := operation.DeleteFileIdsAtOneVolumeServer( + context.Background(), pb.NewServerAddressFromDataNode(target.location.dataNode), vcd.grpcDialOption(), fidList, false) diff --git a/weed/shell/command_volume_fsck.go b/weed/shell/command_volume_fsck.go index 22d96b7b3..ddbdaa3c5 100644 --- a/weed/shell/command_volume_fsck.go +++ b/weed/shell/command_volume_fsck.go @@ -939,7 +939,7 @@ func (c *commandVolumeFsck) purgeFileIdsForOneVolume(volumeId uint32, fileIds [] go func(server pb.ServerAddress, fidList []string) { defer wg.Done() - deleteResults := operation.DeleteFileIdsAtOneVolumeServer(server, c.env.option.GrpcDialOption, fidList, false) + deleteResults := operation.DeleteFileIdsAtOneVolumeServer(context.Background(), server, c.env.option.GrpcDialOption, fidList, false) if deleteResults != nil { resultChan <- deleteResults } diff --git a/weed/stats/metrics.go b/weed/stats/metrics.go index 0c238c711..7b39b943b 100644 --- a/weed/stats/metrics.go +++ b/weed/stats/metrics.go @@ -789,6 +789,14 @@ var ( Help: "Remote-mount object read attempts by source, bucket and cache result. A cold object retried before caching completes records a miss per attempt; paths outside the buckets folder use bucket \"_other\".", }, []string{"source", "bucket", "result"}) + RemoteCacheEvictedCounter = prometheus.NewCounter( + prometheus.CounterOpts{ + Namespace: Namespace, + Subsystem: subsystemRemote, + Name: "cache_evicted_total", + Help: "Remote-mounted objects whose local chunks were evicted under disk pressure.", + }) + UploadErrorCounter = prometheus.NewCounterVec( prometheus.CounterOpts{ Namespace: Namespace, @@ -1072,6 +1080,7 @@ func init() { Gather.MustRegister(S3BucketReadOnlyGauge) Gather.MustRegister(RemoteCacheReadCounter) + Gather.MustRegister(RemoteCacheEvictedCounter) Gather.MustRegister(S3LifecycleDispatchCounter) Gather.MustRegister(S3LifecycleScheduleDepthGauge)