Skip to content
Open
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
91 changes: 91 additions & 0 deletions internal/util/temp_root.go
Original file line number Diff line number Diff line change
@@ -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
}
101 changes: 101 additions & 0 deletions internal/util/temp_root_test.go
Original file line number Diff line number Diff line change
@@ -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())
})
}
43 changes: 43 additions & 0 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,10 @@ limitations under the License.
package main

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

flag "github.com/spf13/pflag"
Expand Down Expand Up @@ -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"
Expand Down Expand Up @@ -160,6 +163,8 @@ func main() {

logger.SetLogger(logger.NewLogger(logOptions))

staleTempDirs := setupTempRoot()

if defaultServiceAccount != "" {
auth.SetDefaultServiceAccount(defaultServiceAccount)
}
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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)
}()
}