From c22b261abd9e7fcaa4df7dc06dbb5038ab9fb9a0 Mon Sep 17 00:00:00 2001 From: Lev Brouk Date: Wed, 20 May 2026 15:50:43 -0700 Subject: [PATCH 1/2] perf(clean-nfs-cache): restore dirfd-relative statx, drain on cancel MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit PR #2256 changed statInDir from fd-relative statx (using an open directory fd) to AT_FDCWD with the full absolute path. On NFS this inflates a single GETATTR into a LOOKUP chain across every path component — 5-10x more RPCs per stat in practice. In prod this showed up as a ~2.5 min first-batch warmup and roughly the first 10 minutes of every hourly run being scan-bound instead of delete-bound. The race #2256 was working around was real: on ctx cancellation, scanDir returned early via ctx.Err(), triggering the deferred df.Close() while Statter goroutines were still calling df.Fd() and unix.Statx on the fd. After Close, the freed fd number can be reused by another goroutine in the same process (a sibling Scanner opening a new directory), so an in-flight Statter can call statx on the wrong fd. With uniform NFS-cache filenames across directories the stat may silently succeed against the wrong file. The right fix is lifetime management, not removing the dependency. This change: - restores statInDir(df *os.File, filename) using fd-relative unix.Statx(int(df.Fd()), filename, ...) on Linux, and matches the signature on macOS (perf there doesn't matter). - changes statReq to carry *os.File again. - in scanDir, counts submitted requests and always drains exactly that many responses from responseCh before returning, even on ctx cancellation. responseCh is buffered to len(filenames) so Statter sends never block, and Statter always emits a response per request it pulls (it doesn't re-check ctx mid-processing). df.Close() can therefore only run after every Statter that holds df is done. Race detector run (go test -race ./cmd/clean-nfs-cache/cleaner/...) passes; existing tests cover the cancel path. --- .../cmd/clean-nfs-cache/cleaner/clean.go | 2 +- .../cmd/clean-nfs-cache/cleaner/scan.go | 49 ++++++++++--------- .../cmd/clean-nfs-cache/cleaner/stat_linux.go | 16 ++++-- .../cmd/clean-nfs-cache/cleaner/stat_osx.go | 4 +- 4 files changed, 40 insertions(+), 31 deletions(-) diff --git a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/clean.go b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/clean.go index 110f4287ba..f7ceb83f71 100644 --- a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/clean.go +++ b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/clean.go @@ -100,7 +100,7 @@ type Candidate struct { } type statReq struct { - dirPath string + df *os.File name string response chan *statReq f *File diff --git a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go index 42ad6dd39d..384aa8e112 100644 --- a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go +++ b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go @@ -60,7 +60,7 @@ func (c *Cleaner) Statter(ctx context.Context, done *sync.WaitGroup) { case <-ctx.Done(): return case req := <-c.statRequestCh: - f, err := c.statInDir(req.dirPath, req.name) + f, err := c.statInDir(req.df, req.name) req.f = f req.err = err req.response <- req @@ -185,7 +185,6 @@ func (c *Cleaner) scanDir(ctx context.Context, path []*Dir) (out *Dir, err error } dirs := make([]*Dir, 0) - nFiles := 0 var filenames []string for _, e := range entries { name := e.Name() @@ -195,46 +194,48 @@ func (c *Cleaner) scanDir(ctx context.Context, path []*Dir) (out *Dir, err error dirs = append(dirs, NewDir(name)) c.DirC.Add(1) } else { - // file - nFiles++ filenames = append(filenames, name) } } - // Submit stat requests using the directory path (not the *os.File). - // The file descriptor df is closed when scanDir returns (defer above), - // but Statter goroutines may still be processing requests concurrently. - // Passing the path avoids a race between df.Close() and df.Fd(). + // Submit stat requests using the directory fd so Statter can use + // fd-relative statx — on NFS this avoids per-component LOOKUP RPCs. + // + // Once a Statter has pulled a request off statRequestCh it will always + // send a response (it does not re-check ctx mid-processing). To make + // the deferred df.Close() safe, we must drain a response for every + // successfully-submitted request before returning, even when ctx is + // canceled mid-loop. responseCh is buffered to len(filenames) so a + // Statter's send back never blocks. responseCh := make(chan *statReq, len(filenames)) + submitted := 0 +submitLoop: for _, name := range filenames { select { case <-ctx.Done(): - return nil, ctx.Err() - case c.statRequestCh <- &statReq{dirPath: absPath, name: name, response: responseCh}: - // submitted + err = ctx.Err() + break submitLoop + case c.statRequestCh <- &statReq{df: df, name: name, response: responseCh}: + submitted++ } } - // get all stat responses - err = nil - files := make([]File, nFiles) - for i := range nFiles { - select { - case <-ctx.Done(): - return nil, ctx.Err() - case resp := <-responseCh: - if resp.err != nil { + files := make([]File, 0, submitted) + for range submitted { + resp := <-responseCh + switch { + case resp.err != nil: + if err == nil { err = resp.err - - continue } - files[i] = *resp.f + case err == nil: + files = append(files, *resp.f) } } if err != nil { return nil, err } - c.FileC.Add(int64(nFiles)) + c.FileC.Add(int64(len(files))) d.mu.Lock() d.Dirs = dirs diff --git a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_linux.go b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_linux.go index 1b1719faed..8754830331 100644 --- a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_linux.go +++ b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_linux.go @@ -4,6 +4,7 @@ package cleaner import ( "fmt" + "os" "path/filepath" "golang.org/x/sys/unix" @@ -30,18 +31,25 @@ func (c *Cleaner) stat(fullPath string) (*Candidate, error) { }, nil } -func (c *Cleaner) statInDir(dirPath string, filename string) (*File, error) { +// statInDir issues an fd-relative statx against an already-open directory fd. +// On NFS this is dramatically cheaper than statx(AT_FDCWD, abs_path), because +// the open dir's file handle is cached on the server and a single GETATTR RPC +// satisfies the call; AT_FDCWD with an absolute path forces a LOOKUP chain +// for every uncached path component. +// +// Safety: the caller must keep df open until this returns. scanDir enforces +// that by draining all in-flight stat responses before its deferred Close. +func (c *Cleaner) statInDir(df *os.File, filename string) (*File, error) { c.StatxC.Add(1) c.StatxInDirC.Add(1) var statx unix.Statx_t - fullPath := filepath.Join(dirPath, filename) - err := unix.Statx(unix.AT_FDCWD, fullPath, + err := unix.Statx(int(df.Fd()), filename, unix.AT_STATX_DONT_SYNC|unix.AT_SYMLINK_NOFOLLOW|unix.AT_NO_AUTOMOUNT, unix.STATX_ATIME|unix.STATX_SIZE, &statx, ) if err != nil { - return nil, fmt.Errorf("failed to statx %q: %w", fullPath, err) + return nil, fmt.Errorf("failed to statx %q: %w", filepath.Join(df.Name(), filename), err) } return &File{ diff --git a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_osx.go b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_osx.go index 2b793f4397..ce18be2907 100644 --- a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_osx.go +++ b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_osx.go @@ -31,10 +31,10 @@ func (c *Cleaner) stat(path string) (*Candidate, error) { }, nil } -func (c *Cleaner) statInDir(dirPath string, filename string) (*File, error) { +func (c *Cleaner) statInDir(df *os.File, filename string) (*File, error) { c.StatxInDirC.Add(1) // performance on OS X does not matter, so just use the full stat - cand, err := c.stat(filepath.Join(dirPath, filename)) + cand, err := c.stat(filepath.Join(df.Name(), filename)) if err != nil { return nil, err } From 690e3e40730e6e1c3594c46bf9be0b655ad61c96 Mon Sep 17 00:00:00 2001 From: Lev Brouk Date: Wed, 20 May 2026 16:10:02 -0700 Subject: [PATCH 2/2] fix(clean-nfs-cache): truthful shutdown logs MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two cosmetic issues both visible at the end of every successful run: 1. When preRun computes that current disk usage is already at or below --disk-usage-target-percent, TargetBytesToDelete is left at 0. The main loop's first check (DeletedBytes(0) >= TargetBytesToDelete(0)) then immediately logged "target bytes deleted reached, draining remaining candidates" — implying the cleaner had done work, when in fact there was nothing to do. Now short-circuited in main.go with an explicit "disk already at or below target, nothing to do". 2. Each Scanner goroutine in flight at drain time logged "error during scanning: context canceled" before exiting cleanly, making routine shutdown look like a failure. Now suppressed: a context.Canceled / context.DeadlineExceeded from FindCandidate is treated as a shutdown signal and the Scanner returns without logging or writing to errCh. Also fixed a pre-existing typo in the scanner error log field name ("continousCount" -> "continuousCount"). --- .../orchestrator/cmd/clean-nfs-cache/cleaner/scan.go | 8 +++++++- packages/orchestrator/cmd/clean-nfs-cache/main.go | 11 +++++++++++ 2 files changed, 18 insertions(+), 1 deletion(-) diff --git a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go index 384aa8e112..5b517e4a05 100644 --- a/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go +++ b/packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go @@ -35,10 +35,15 @@ func (c *Cleaner) Scanner(ctx context.Context, candidateCh chan<- *Candidate, er continue + case errors.Is(err, context.Canceled), errors.Is(err, context.DeadlineExceeded): + // Shutdown in progress; the outer select will exit on the + // next iteration. Don't log it as an error or pollute errCh. + return + default: if !errors.Is(err, ErrNoFiles) { c.Info(ctx, "error during scanning", - zap.Int("continousCount", continuousErrors), + zap.Int("continuousCount", continuousErrors), zap.Error(err)) } continuousErrors++ @@ -214,6 +219,7 @@ submitLoop: select { case <-ctx.Done(): err = ctx.Err() + break submitLoop case c.statRequestCh <- &statReq{df: df, name: name, response: responseCh}: submitted++ diff --git a/packages/orchestrator/cmd/clean-nfs-cache/main.go b/packages/orchestrator/cmd/clean-nfs-cache/main.go index b03aaaef15..885bbfac77 100644 --- a/packages/orchestrator/cmd/clean-nfs-cache/main.go +++ b/packages/orchestrator/cmd/clean-nfs-cache/main.go @@ -70,6 +70,17 @@ func main() { zap.Int("max_concurrent_delete", opts.MaxConcurrentDelete), ) + // preRun leaves both targets at 0 when disk usage is already at or + // below disk-usage-target-percent. Short-circuit here so we don't spin + // up workers just to immediately drain (which also produced a misleading + // "target bytes deleted reached" log). + if opts.TargetBytesToDelete == 0 && opts.TargetFilesToDelete == 0 { + log.Info(ctx, "disk already at or below target, nothing to do", + zap.Float64("target_disk_usage_percent", opts.TargetDiskUsagePercent)) + + return + } + c := cleaner.NewCleaner(opts, log) if err = c.Clean(ctx); err != nil { return