diff --git a/pkg/state/issue_768_test.go b/pkg/state/issue_768_test.go new file mode 100644 index 00000000..957707c5 --- /dev/null +++ b/pkg/state/issue_768_test.go @@ -0,0 +1,249 @@ +package state + +import ( + "fmt" + "os" + "path/filepath" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.uber.org/zap" + + "github.com/helmfile/helmfile/pkg/exectest" + "github.com/helmfile/helmfile/pkg/filesystem" +) + +// mockFetchHelm wraps exectest.Helm to count Fetch calls and create +// a fake chart directory so findChartDirectory succeeds. +type mockFetchHelm struct { + *exectest.Helm + fetchCount atomic.Int32 +} + +func (m *mockFetchHelm) Fetch(chart string, flags ...string) error { + // Extract the --untardir path from flags and create a fake Chart.yaml + untarDir := "" + for i, f := range flags { + if f == "--untardir" && i+1 < len(flags) { + untarDir = flags[i+1] + break + } + } + if untarDir != "" { + chartDir := filepath.Join(untarDir, chart) + _ = os.MkdirAll(chartDir, 0755) + chartYaml := filepath.Join(chartDir, "Chart.yaml") + _ = os.WriteFile(chartYaml, []byte("apiVersion: v2\nname: test\nversion: 1.0.0\n"), 0644) + } + m.fetchCount.Add(1) + // Simulate download time to widen the race window + time.Sleep(50 * time.Millisecond) + return nil +} + +// TestForcedDownloadChartSerializesSameChart verifies that concurrent calls +// to forcedDownloadChart with the same chart+version but different release +// names are serialized so helm.Fetch is called exactly once. +// This is the core fix for issue #768. +func TestForcedDownloadChartSerializesSameChart(t *testing.T) { + logger, err := zap.NewDevelopment() + require.NoError(t, err) + + tempDir := t.TempDir() + + // Use unique chart name to avoid polluting global cache across tests + chartName := fmt.Sprintf("testrepo-768/unique-chart-%d", time.Now().UnixNano()) + + st := &HelmState{ + logger: logger.Sugar(), + fs: filesystem.DefaultFileSystem(), + basePath: tempDir, + } + + mockHelm := &mockFetchHelm{Helm: &exectest.Helm{}} + + numReleases := 5 + var wg sync.WaitGroup + + for i := 0; i < numReleases; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + release := &ReleaseSpec{ + Name: fmt.Sprintf("release-%d", idx), + Chart: chartName, + Version: "1.0.0", + } + opts := ChartPrepareOptions{ + SkipRefresh: true, + SkipDeps: true, + } + _, err := st.forcedDownloadChart(chartName, tempDir, release, mockHelm, opts) + assert.NoError(t, err) + }(i) + } + + wg.Wait() + + // helm.Fetch should be called exactly once despite 5 concurrent releases + assert.Equal(t, int32(1), mockHelm.fetchCount.Load(), + "helm.Fetch must be called exactly once for concurrent releases with the same chart+version") +} + +// TestForcedDownloadChartAllowsDifferentCharts verifies that downloads of +// different charts are NOT serialized against each other (only same-chart +// downloads are serialized). +func TestForcedDownloadChartAllowsDifferentCharts(t *testing.T) { + logger, err := zap.NewDevelopment() + require.NoError(t, err) + + tempDir := t.TempDir() + ts := time.Now().UnixNano() + + st := &HelmState{ + logger: logger.Sugar(), + fs: filesystem.DefaultFileSystem(), + basePath: tempDir, + } + + mockHelm := &mockFetchHelm{Helm: &exectest.Helm{}} + + // Two different charts + charts := []string{ + fmt.Sprintf("testrepo-768/chart-a-%d", ts), + fmt.Sprintf("testrepo-768/chart-b-%d", ts), + } + + var wg sync.WaitGroup + for _, chartName := range charts { + wg.Add(1) + go func(cn string) { + defer wg.Done() + release := &ReleaseSpec{ + Name: "release-" + cn, + Chart: cn, + Version: "1.0.0", + } + opts := ChartPrepareOptions{ + SkipRefresh: true, + SkipDeps: true, + } + _, e := st.forcedDownloadChart(cn, tempDir, release, mockHelm, opts) + assert.NoError(t, e) + }(chartName) + } + + wg.Wait() + + // Both different charts should be fetched + assert.Equal(t, int32(2), mockHelm.fetchCount.Load(), + "different charts should not block each other") +} + +// TestWithChartOperationLockSerializesSameChart verifies that +// withChartOperationLock serializes concurrent helm operations for +// releases using the same remote chart (issue #768). +func TestWithChartOperationLockSerializesSameChart(t *testing.T) { + logger, err := zap.NewDevelopment() + require.NoError(t, err) + + ts := time.Now().UnixNano() + chartName := fmt.Sprintf("testrepo-768/op-chart-%d", ts) + + st := &HelmState{ + logger: logger.Sugar(), + } + + var concurrentCount atomic.Int32 + var maxConcurrent atomic.Int32 + + numReleases := 5 + var wg sync.WaitGroup + + for i := 0; i < numReleases; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + release := &ReleaseSpec{ + Name: fmt.Sprintf("release-%d", idx), + Chart: chartName, + Version: "1.0.0", + // ChartPath empty => remote chart => lock acquired + } + _ = st.withChartOperationLock(release, chartName, func() error { + cur := concurrentCount.Add(1) + for { + old := maxConcurrent.Load() + if cur <= old || maxConcurrent.CompareAndSwap(old, cur) { + break + } + } + time.Sleep(20 * time.Millisecond) + concurrentCount.Add(-1) + return nil + }) + }(i) + } + + wg.Wait() + + // Same-chart operations must be serialized: max concurrency should be 1 + assert.Equal(t, int32(1), maxConcurrent.Load(), + "concurrent helm operations for the same chart must be serialized") +} + +// TestWithChartOperationLockNoLockForLocalChart verifies that +// withChartOperationLock does NOT serialize operations when ChartPath +// is set (local or pre-fetched charts). +func TestWithChartOperationLockNoLockForLocalChart(t *testing.T) { + logger, err := zap.NewDevelopment() + require.NoError(t, err) + + ts := time.Now().UnixNano() + chartName := fmt.Sprintf("testrepo-768/local-chart-%d", ts) + + st := &HelmState{ + logger: logger.Sugar(), + } + + var concurrentCount atomic.Int32 + var maxConcurrent atomic.Int32 + + numReleases := 5 + var wg sync.WaitGroup + + for i := 0; i < numReleases; i++ { + wg.Add(1) + go func(idx int) { + defer wg.Done() + release := &ReleaseSpec{ + Name: fmt.Sprintf("release-%d", idx), + Chart: chartName, + Version: "1.0.0", + ChartPath: "/some/local/path", // ChartPath set => no lock + } + _ = st.withChartOperationLock(release, chartName, func() error { + cur := concurrentCount.Add(1) + for { + old := maxConcurrent.Load() + if cur <= old || maxConcurrent.CompareAndSwap(old, cur) { + break + } + } + time.Sleep(20 * time.Millisecond) + concurrentCount.Add(-1) + return nil + }) + }(i) + } + + wg.Wait() + + // Local charts should NOT be serialized: all should run concurrently + assert.Greater(t, maxConcurrent.Load(), int32(1), + "local/pre-fetched charts should not be serialized") +} diff --git a/pkg/state/state.go b/pkg/state/state.go index 0e118ee9..4aba56bd 100644 --- a/pkg/state/state.go +++ b/pkg/state/state.go @@ -1237,7 +1237,9 @@ func (st *HelmState) SyncReleases(affectedReleases *AffectedReleases, helm helme // active. When tracking isn't running it's the original // helm — either path is safe to call directly. releaseHelm := trackHandle.Helm - helmErr := releaseHelm.SyncRelease(context, release.Name, chart, release.Namespace, flags...) + helmErr := st.withChartOperationLock(release, chart, func() error { + return releaseHelm.SyncRelease(context, release.Name, chart, release.Namespace, flags...) + }) if helmErr != nil && trackStarted && trackHandle.WasHelmKilled() { // The kubedog safety valve sent SIGINT to helm because // the cluster confirmed convergence while helm was @@ -1336,7 +1338,9 @@ func (st *HelmState) SyncReleases(affectedReleases *AffectedReleases, helm helme } func (st *HelmState) performSyncOrReinstallOfRelease(affectedReleases *AffectedReleases, helm helmexec.Interface, context helmexec.HelmContext, release *ReleaseSpec, chart string, m *sync.Mutex, flags ...string) *ReleaseError { - if err := helm.SyncRelease(context, release.Name, chart, release.Namespace, flags...); err != nil { + if err := st.withChartOperationLock(release, chart, func() error { + return helm.SyncRelease(context, release.Name, chart, release.Namespace, flags...) + }); err != nil { st.logger.Debugf("update strategy - sync failed: %s", err.Error()) // Only fail if a different error than forbidden updates if !strings.Contains(err.Error(), "Forbidden: updates") { @@ -1387,7 +1391,9 @@ func (st *HelmState) performSyncOrReinstallOfRelease(affectedReleases *AffectedR } m.Unlock() } - if err := helm.SyncRelease(context, release.Name, chart, release.Namespace, flags...); err != nil { + if err := st.withChartOperationLock(release, chart, func() error { + return helm.SyncRelease(context, release.Name, chart, release.Namespace, flags...) + }); err != nil { m.Lock() affectedReleases.Failed = append(affectedReleases.Failed, release) m.Unlock() @@ -1552,6 +1558,33 @@ func (st *HelmState) getNamedRWMutex(name string) *sync.RWMutex { return actualMu.(*sync.RWMutex) } +// withChartOperationLock serializes helm operations (upgrade, diff) per +// chart+version for remote charts to prevent concurrent download races on +// helm's repository cache (issue #768). When multiple releases use the same +// remote chart, concurrent `helm upgrade`/`helm diff` calls race on helm's +// internal cache file rename, causing "Access is denied" errors on Windows. +// +// For local or pre-fetched charts (release.ChartPath is set), no lock is +// acquired because the chart is already available locally and helm won't +// download it. +// +// Trade-off: the lock wraps the entire helm operation (download + template + +// deploy), not just the download. Same-chart releases are fully serialized. +// This is unavoidable without pre-fetching because helm downloads and deploys +// atomically. Releases with different charts are unaffected and remain fully +// parallel. If kubedog tracking is enabled, note that it starts before this +// lock is acquired, so a release waiting on the lock may have its kubedog +// watcher active before helm begins. +func (st *HelmState) withChartOperationLock(release *ReleaseSpec, chart string, fn func() error) error { + if release.ChartPath != "" { + return fn() + } + mu := st.getNamedRWMutex("helm-op:" + chart + ":" + release.Version) + mu.Lock() + defer mu.Unlock() + return fn() +} + type PrepareChartKey struct { Namespace, Name, KubeContext string } @@ -1947,17 +1980,37 @@ func (st *HelmState) processLocalChart(normalizedChart, dir string, release *Rel // forcedDownloadChart handles forced chart downloads. // Locks are acquired during download and released immediately after. +// A per-chart+version mutex serializes downloads within the process so that +// concurrent releases using the same chart don't race on helm's repository +// cache (issue #768). func (st *HelmState) forcedDownloadChart(chartName, dir string, release *ReleaseSpec, helm helmexec.Interface, opts ChartPrepareOptions) (string, error) { - // Check global chart cache first for non-OCI charts + cacheKey := st.getChartCacheKey(release) + + // Fast path: check in-process cache without acquiring any lock. // If found, another worker in this process already downloaded the chart. // We don't need to acquire a lock - the tempDir won't be deleted until // after helm operations complete (cleanup is deferred in withPreparedCharts). - cacheKey := st.getChartCacheKey(release) if cachedPath, exists := st.checkChartCache(cacheKey); exists && st.fs.DirectoryExistsAt(cachedPath) { st.logger.Debugf("Chart %s:%s already downloaded, using cached version at %s", chartName, release.Version, cachedPath) return cachedPath, nil } + // Serialize downloads per chart+version within this process. + // The chartPath generated below includes the release name (see + // DefaultFetchOutputDirTemplate), so per-chartPath file locks do NOT + // prevent concurrent `helm fetch` of the same chart by different releases. + // Without this mutex, two workers downloading the same chart concurrently + // race on helm's repository cache on Windows (issue #768). + downloadMu := st.getNamedRWMutex("chart-download:" + cacheKey.Chart + ":" + cacheKey.Version) + downloadMu.Lock() + defer downloadMu.Unlock() + + // Double-check: another worker may have completed the download while we waited. + if cachedPath, exists := st.checkChartCache(cacheKey); exists && st.fs.DirectoryExistsAt(cachedPath) { + st.logger.Debugf("Chart %s:%s downloaded by another worker, using cached version at %s", chartName, release.Version, cachedPath) + return cachedPath, nil + } + chartPath, err := st.generateChartPath(chartName, dir, release, opts.OutputDirTemplate) if err != nil { return "", err @@ -3056,7 +3109,9 @@ func (st *HelmState) DiffReleases(helm helmexec.Interface, additionalValues []st if prep.upgradeDueToSkippedDiff { results <- diffResult{release, &ReleaseError{ReleaseSpec: release, err: nil, Code: HelmDiffExitCodeChanged}, buf} - } else if err := helm.DiffRelease(st.createHelmContextWithWriter(release, buf), release.Name, chartPath, release.Namespace, releaseSuppressDiff, flags...); err != nil { + } else if err := st.withChartOperationLock(release, chartPath, func() error { + return helm.DiffRelease(st.createHelmContextWithWriter(release, buf), release.Name, chartPath, release.Namespace, releaseSuppressDiff, flags...) + }); err != nil { switch e := err.(type) { case helmexec.ExitError: // Propagate any non-zero exit status from the external command like `helm` that is failed under the hood @@ -5622,6 +5677,8 @@ func (st *HelmState) acquireExclusiveLock(result *chartLockResult, chartPath str // getOCIChart downloads or retrieves an OCI chart from cache. // Locks are acquired during download and released immediately after. +// A per-chart+version mutex serializes downloads within the process so that +// concurrent releases using the same OCI chart don't race (issue #768). func (st *HelmState) getOCIChart(release *ReleaseSpec, tempDir string, helm helmexec.Interface, opts ChartPrepareOptions) (*string, error) { qualifiedChartName, chartName, chartVersion, err := st.getOCIQualifiedChartName(release) if err != nil { @@ -5632,16 +5689,25 @@ func (st *HelmState) getOCIChart(release *ReleaseSpec, tempDir string, helm helm return nil, nil } - // Check global chart cache first (in-memory cache) - // If found, another worker in this process already downloaded the chart. - // We don't need to acquire a lock - the tempDir won't be deleted until - // after helm operations complete (cleanup is deferred in withPreparedCharts). cacheKey := st.getChartCacheKey(release) + + // Fast path: check in-process cache without acquiring any lock. if cachedPath, exists := st.checkChartCache(cacheKey); exists && st.fs.DirectoryExistsAt(cachedPath) { st.logger.Debugf("OCI chart %s:%s already downloaded, using cached version at %s", chartName, chartVersion, cachedPath) return &cachedPath, nil } + // Serialize downloads per chart+version within this process (issue #768). + downloadMu := st.getNamedRWMutex("chart-download:" + cacheKey.Chart + ":" + cacheKey.Version) + downloadMu.Lock() + defer downloadMu.Unlock() + + // Double-check: another worker may have completed the download while we waited. + if cachedPath, exists := st.checkChartCache(cacheKey); exists && st.fs.DirectoryExistsAt(cachedPath) { + st.logger.Debugf("OCI chart %s:%s downloaded by another worker, using cached version at %s", chartName, chartVersion, cachedPath) + return &cachedPath, nil + } + if opts.OutputDirTemplate == "" { tempDir = remote.CacheDir() }