From 91484bdcc228b27b98c3dad9488723d968088146 Mon Sep 17 00:00:00 2001 From: Andrii Bryndzak Date: Wed, 30 Sep 2026 19:32:00 +0300 Subject: [PATCH] fix: purge stale tmp dirs at startup When the controller exits without running its deferred cleanup, after losing its leader lease or being OOM killed, the temporary directories of the reconciles in flight stay behind. They survive container restarts on an emptyDir volume, so a controller that restarts often fills the node. The controller now keeps its temporary files in a directory of its own under the system temp dir and points TMPDIR at it, so the files created by the Helm and Git libraries land there too. At startup, whatever a previous process left in that directory is moved aside and removed in the background, the same way kustomize-controller purges its tmp dirs. Signed-off-by: Andrii Bryndzak --- internal/util/temp_root.go | 91 ++++++++++++++++++++++++++++ internal/util/temp_root_test.go | 101 ++++++++++++++++++++++++++++++++ main.go | 43 ++++++++++++++ 3 files changed, 235 insertions(+) create mode 100644 internal/util/temp_root.go create mode 100644 internal/util/temp_root_test.go diff --git a/internal/util/temp_root.go b/internal/util/temp_root.go new file mode 100644 index 000000000..d56523797 --- /dev/null +++ b/internal/util/temp_root.go @@ -0,0 +1,91 @@ +/* +Copyright 2026 The Flux authors + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package util + +import ( + "context" + "errors" + "fmt" + "io/fs" + "os" + "path/filepath" + "strings" + "time" + + "github.com/go-logr/logr" +) + +// staleSuffix is appended to a temp root that a previous process left behind, +// when it is moved aside to be purged. +const staleSuffix = ".stale-" + +// PrepareTempRoot gives the process an empty directory at root for its +// temporary files. Anything a previous process left in root, because it +// exited without running its deferred cleanup, is moved aside first. It +// returns the paths that hold leftovers, including those a previous purge +// did not finish, so that the caller can remove them with PurgeTempRoots. +// It must be called before the reconcilers start. +func PrepareTempRoot(root string) ([]string, error) { + _, err := os.Lstat(root) + switch { + case err == nil: + stale := fmt.Sprintf("%s%s%d", root, staleSuffix, time.Now().UnixNano()) + if err := os.Rename(root, stale); err != nil { + return nil, fmt.Errorf("failed to move aside %s: %w", root, err) + } + case !errors.Is(err, fs.ErrNotExist): + return nil, fmt.Errorf("failed to stat %s: %w", root, err) + } + + if err := os.MkdirAll(root, 0o700); err != nil { + return nil, fmt.Errorf("failed to create %s: %w", root, err) + } + + parent, base := filepath.Split(root) + entries, err := os.ReadDir(parent) + if err != nil { + return nil, fmt.Errorf("failed to list %s: %w", parent, err) + } + var stale []string + for _, entry := range entries { + if strings.HasPrefix(entry.Name(), base+staleSuffix) { + stale = append(stale, filepath.Join(parent, entry.Name())) + } + } + return stale, nil +} + +// PurgeTempRoots removes the given paths, logging every removal at info +// level and every failure at error level. It stops as soon as ctx is done +// and returns the number of paths that were removed. +func PurgeTempRoots(ctx context.Context, log logr.Logger, paths []string) int { + purged := 0 + for i, path := range paths { + if err := ctx.Err(); err != nil { + log.Error(err, "aborted purge of stale tmp dirs", + "purged", purged, "remaining", len(paths)-i) + return purged + } + if err := os.RemoveAll(path); err != nil { + log.Error(err, "failed to remove stale tmp dir", "path", path) + continue + } + log.Info("removed stale tmp dir", "path", path) + purged++ + } + return purged +} diff --git a/internal/util/temp_root_test.go b/internal/util/temp_root_test.go new file mode 100644 index 000000000..6dd1cf44e --- /dev/null +++ b/internal/util/temp_root_test.go @@ -0,0 +1,101 @@ +/* +Copyright 2026 The Flux authors + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package util + +import ( + "context" + "os" + "path/filepath" + "testing" + + "github.com/go-logr/logr" + . "github.com/onsi/gomega" +) + +func TestPrepareTempRoot(t *testing.T) { + t.Run("creates the root when it does not exist", func(t *testing.T) { + g := NewWithT(t) + root := filepath.Join(t.TempDir(), "source-controller") + + stale, err := PrepareTempRoot(root) + g.Expect(err).ToNot(HaveOccurred()) + g.Expect(stale).To(BeEmpty()) + g.Expect(root).To(BeADirectory()) + }) + + t.Run("moves a previous process's leftovers out of the root", func(t *testing.T) { + g := NewWithT(t) + root := filepath.Join(t.TempDir(), "source-controller") + leftover := filepath.Join(root, "gitrepository-default-foo-123", "repo") + g.Expect(os.MkdirAll(leftover, 0o700)).To(Succeed()) + g.Expect(os.WriteFile(filepath.Join(root, "chart-index-1.yaml"), []byte("x"), 0o600)).To(Succeed()) + + stale, err := PrepareTempRoot(root) + g.Expect(err).ToNot(HaveOccurred()) + g.Expect(stale).To(HaveLen(1)) + g.Expect(filepath.Join(stale[0], "gitrepository-default-foo-123", "repo")).To(BeADirectory()) + g.Expect(filepath.Join(stale[0], "chart-index-1.yaml")).To(BeARegularFile()) + + entries, err := os.ReadDir(root) + g.Expect(err).ToNot(HaveOccurred()) + g.Expect(entries).To(BeEmpty()) + }) + + t.Run("lists leftovers a previous purge did not finish", func(t *testing.T) { + g := NewWithT(t) + parent := t.TempDir() + root := filepath.Join(parent, "source-controller") + unfinished := filepath.Join(parent, "source-controller"+staleSuffix+"1") + unrelated := filepath.Join(parent, "other-controller"+staleSuffix+"1") + for _, dir := range []string{root, unfinished, unrelated} { + g.Expect(os.MkdirAll(dir, 0o700)).To(Succeed()) + } + + stale, err := PrepareTempRoot(root) + g.Expect(err).ToNot(HaveOccurred()) + g.Expect(stale).To(ContainElement(unfinished)) + g.Expect(stale).ToNot(ContainElement(unrelated)) + g.Expect(stale).To(HaveLen(2)) + }) +} + +func TestPurgeTempRoots(t *testing.T) { + t.Run("removes every path", func(t *testing.T) { + g := NewWithT(t) + parent := t.TempDir() + a := filepath.Join(parent, "a", "nested") + b := filepath.Join(parent, "b") + g.Expect(os.MkdirAll(a, 0o700)).To(Succeed()) + g.Expect(os.MkdirAll(b, 0o700)).To(Succeed()) + + purged := PurgeTempRoots(context.Background(), logr.Discard(), []string{filepath.Join(parent, "a"), b}) + g.Expect(purged).To(Equal(2)) + g.Expect(filepath.Join(parent, "a")).ToNot(BeADirectory()) + g.Expect(b).ToNot(BeADirectory()) + }) + + t.Run("stops when the context is done", func(t *testing.T) { + g := NewWithT(t) + dir := t.TempDir() + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + purged := PurgeTempRoots(ctx, logr.Discard(), []string{dir}) + g.Expect(purged).To(Equal(0)) + g.Expect(dir).To(BeADirectory()) + }) +} diff --git a/main.go b/main.go index 75d897bd8..ef432f7b8 100644 --- a/main.go +++ b/main.go @@ -17,8 +17,10 @@ limitations under the License. package main import ( + "context" "fmt" "os" + "path/filepath" "time" flag "github.com/spf13/pflag" @@ -64,6 +66,7 @@ import ( "github.com/fluxcd/source-controller/internal/features" "github.com/fluxcd/source-controller/internal/helm" scosign "github.com/fluxcd/source-controller/internal/oci/cosign" + "github.com/fluxcd/source-controller/internal/util" ) const controllerName = "source-controller" @@ -160,6 +163,8 @@ func main() { logger.SetLogger(logger.NewLogger(logOptions)) + staleTempDirs := setupTempRoot() + if defaultServiceAccount != "" { auth.SetDefaultServiceAccount(defaultServiceAccount) } @@ -321,6 +326,8 @@ func main() { } }() + purgeStaleTempDirs(ctx, staleTempDirs) + setupLog.Info("starting manager") if err := mgr.Start(ctx); err != nil { setupLog.Error(err, "problem running manager") @@ -450,3 +457,39 @@ func envOrDefault(envName, defaultValue string) string { return defaultValue } + +// setupTempRoot points TMPDIR at a directory of this controller's own, so +// that every temporary file and directory it creates, including those of the +// Helm and Git libraries, lands in one place. Whatever a previous process +// left there, because it exited without running its deferred cleanup, is +// moved aside and returned for purgeStaleTempDirs to remove. +func setupTempRoot() []string { + root := filepath.Join(os.TempDir(), controllerName) + stale, err := util.PrepareTempRoot(root) + if err != nil { + setupLog.Error(err, "unable to set up tmp dir") + os.Exit(1) + } + if err := os.Setenv("TMPDIR", root); err != nil { + setupLog.Error(err, "unable to set TMPDIR") + os.Exit(1) + } + return stale +} + +// purgeStaleTempDirs removes the leftovers found by setupTempRoot in the +// background, so that the manager startup is not delayed. +func purgeStaleTempDirs(ctx context.Context, dirs []string) { + const purgeTimeout = 2 * time.Minute + + if len(dirs) == 0 { + return + } + + setupLog.Info("purging stale tmp dirs", "count", len(dirs)) + go func() { + ctx, cancel := context.WithTimeout(ctx, purgeTimeout) + defer cancel() + util.PurgeTempRoots(ctx, setupLog, dirs) + }() +}