mirror of
https://github.com/helmfile/helmfile.git
synced 2026-09-30 06:44:35 +02:00
When multiple releases in a helmfile use the same remote chart, concurrent helm upgrade/diff calls race on helm's internal repository cache file rename, causing intermittent 'cannot rename: Access is denied' errors on Windows. This fix introduces two layers of serialization: 1. withChartOperationLock (operation-level): wraps SyncRelease/DiffRelease calls with a per-chart+version mutex. Only applies to remote charts (release.ChartPath is empty); local/pre-fetched/OCI charts bypass the lock entirely. Different charts remain fully parallel. 2. Per-chart+version download mutex in forcedDownloadChart/getOCIChart (download-level): uses double-check locking to ensure only one helm fetch runs per unique chart+version within a process. The fix does NOT change chart paths passed to helm, preserving backward compatibility with all existing behavior and tests. Trade-off: same-chart releases are fully serialized (the entire helm operation including deployment, not just download). This is unavoidable without pre-fetching because helm downloads and deploys atomically. Releases with different charts are unaffected. Fixes #768 Signed-off-by: yxxhero <aiopsclub@163.com>
This commit is contained in:
@@ -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")
|
||||
}
|
||||
+76
-10
@@ -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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user