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) + }() +}