Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,7 @@ type Candidate struct {
}

type statReq struct {
dirPath string
df *os.File
name string
response chan *statReq
f *File
Expand Down
57 changes: 32 additions & 25 deletions packages/orchestrator/cmd/clean-nfs-cache/cleaner/scan.go
Original file line number Diff line number Diff line change
Expand Up @@ -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++
Expand All @@ -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
Expand Down Expand Up @@ -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()
Expand All @@ -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
Expand Down
16 changes: 12 additions & 4 deletions packages/orchestrator/cmd/clean-nfs-cache/cleaner/stat_linux.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ package cleaner

import (
"fmt"
"os"
"path/filepath"

"golang.org/x/sys/unix"
Expand All @@ -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{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
11 changes: 11 additions & 0 deletions packages/orchestrator/cmd/clean-nfs-cache/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Guard early exit with disk-usage target check

This early return now bypasses Cleaner.validateOptions() for configurations where both byte/file targets are zero but disk-usage-target-percent is also zero (for example --disk-usage-target-percent=0 with no explicit delete targets), so the command exits successfully with “nothing to do” instead of surfacing the existing usage error. That turns an invalid configuration into a silent no-op and can hide misconfigured automation.

Useful? React with 👍 / 👎.

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
Expand Down
Loading