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..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++ @@ -60,7 +65,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 +190,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 +199,49 @@ 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 } 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