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)