fix: configure kubedog rate limiter to prevent context cancellation (#2446)

This commit is contained in:
yxxhero
2026-03-03 19:24:28 +08:00
committed by GitHub
parent 6e21671228
commit ce09f560d9
11 changed files with 887 additions and 48 deletions
+14
View File
@@ -18,12 +18,16 @@ type TrackOptions struct {
Logs bool
LogsSince time.Duration
Filter *resource.FilterConfig
QPS float32
Burst int
}
func NewTrackOptions() *TrackOptions {
return &TrackOptions{
Timeout: 5 * time.Minute,
LogsSince: 10 * time.Minute,
QPS: 100,
Burst: 200,
}
}
@@ -41,3 +45,13 @@ func (o *TrackOptions) WithFilterConfig(config *resource.FilterConfig) *TrackOpt
o.Filter = config
return o
}
func (o *TrackOptions) WithQPS(qps float32) *TrackOptions {
o.QPS = qps
return o
}
func (o *TrackOptions) WithBurst(burst int) *TrackOptions {
o.Burst = burst
return o
}
+133 -38
View File
@@ -3,16 +3,23 @@ package kubedog
import (
"context"
"fmt"
"math"
"os"
"strings"
"sync"
"time"
"github.com/werf/kubedog/pkg/kube"
"github.com/werf/kubedog/pkg/tracker"
"github.com/werf/kubedog/pkg/trackers/rollout/multitrack"
"github.com/werf/kubedog-for-werf-helm/pkg/tracker"
"github.com/werf/kubedog-for-werf-helm/pkg/trackers/rollout/multitrack"
"go.uber.org/zap"
"k8s.io/apimachinery/pkg/api/meta"
"k8s.io/client-go/discovery"
"k8s.io/client-go/discovery/cached/memory"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/restmapper"
"k8s.io/client-go/tools/clientcmd"
"github.com/helmfile/helmfile/pkg/resource"
)
@@ -20,19 +27,32 @@ import (
type cacheKey struct {
kubeContext string
kubeconfig string
qps float32
burst int
}
type clientCacheEntry struct {
clientSet kubernetes.Interface
dynamicClient dynamic.Interface
restConfig *rest.Config
discovery discovery.CachedDiscoveryInterface
mapper meta.RESTMapper
}
var (
kubeInitMu sync.Mutex
clientCache = make(map[cacheKey]kubernetes.Interface)
clientCache = make(map[cacheKey]clientCacheEntry)
)
type Tracker struct {
logger *zap.SugaredLogger
clientSet kubernetes.Interface
trackOptions *TrackOptions
filter *resource.ResourceFilter
namespace string
logger *zap.SugaredLogger
clientSet kubernetes.Interface
dynamicClient dynamic.Interface
discovery discovery.CachedDiscoveryInterface
mapper meta.RESTMapper
trackOptions *TrackOptions
filter *resource.ResourceFilter
namespace string
}
type TrackerConfig struct {
@@ -41,6 +61,8 @@ type TrackerConfig struct {
KubeContext string
Kubeconfig string
TrackOptions *TrackOptions
KubedogQPS *float32
KubedogBurst *int
}
func NewTracker(config *TrackerConfig) (*Tracker, error) {
@@ -54,58 +76,115 @@ func NewTracker(config *TrackerConfig) (*Tracker, error) {
kubeconfig = os.Getenv("KUBECONFIG")
}
clientSet, err := getOrCreateClient(config.KubeContext, kubeconfig)
if err != nil {
return nil, fmt.Errorf("failed to initialize kubernetes client: %w", err)
}
options := config.TrackOptions
if options == nil {
options = NewTrackOptions()
}
qps := options.QPS
if config.KubedogQPS != nil {
qps = *config.KubedogQPS
}
burst := options.Burst
if config.KubedogBurst != nil {
burst = *config.KubedogBurst
}
if qps <= 0 || math.IsInf(float64(qps), 0) || math.IsNaN(float64(qps)) {
return nil, fmt.Errorf("invalid kubedog QPS %v: must be > 0 and finite", qps)
}
if burst < 1 {
return nil, fmt.Errorf("invalid kubedog burst %v: must be >= 1", burst)
}
cacheEntry, err := getOrCreateClients(config.KubeContext, kubeconfig, qps, burst)
if err != nil {
return nil, fmt.Errorf("failed to initialize kubernetes clients: %w", err)
}
var filter *resource.ResourceFilter
if options.Filter != nil {
filter = resource.NewResourceFilter(options.Filter, logger)
}
return &Tracker{
logger: logger,
clientSet: clientSet,
trackOptions: options,
filter: filter,
namespace: config.Namespace,
logger: logger,
clientSet: cacheEntry.clientSet,
dynamicClient: cacheEntry.dynamicClient,
discovery: cacheEntry.discovery,
mapper: cacheEntry.mapper,
trackOptions: options,
filter: filter,
namespace: config.Namespace,
}, nil
}
func getOrCreateClient(kubeContext, kubeconfig string) (kubernetes.Interface, error) {
func getOrCreateClients(kubeContext, kubeconfig string, qps float32, burst int) (clientCacheEntry, error) {
key := cacheKey{
kubeContext: kubeContext,
kubeconfig: kubeconfig,
qps: qps,
burst: burst,
}
kubeInitMu.Lock()
if cache, ok := clientCache[key]; ok {
kubeInitMu.Unlock()
return cache, nil
}
kubeInitMu.Unlock()
loadingRules := clientcmd.NewDefaultClientConfigLoadingRules()
if kubeconfig != "" {
loadingRules.ExplicitPath = kubeconfig
}
overrides := &clientcmd.ConfigOverrides{}
if kubeContext != "" {
overrides.CurrentContext = kubeContext
}
cc := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, overrides)
restConfig, err := cc.ClientConfig()
if err != nil {
return clientCacheEntry{}, fmt.Errorf("failed to load kubeconfig: %w", err)
}
restConfig.QPS = qps
restConfig.Burst = burst
clientSet, err := kubernetes.NewForConfig(restConfig)
if err != nil {
return clientCacheEntry{}, fmt.Errorf("failed to create kubernetes client: %w", err)
}
dynamicClient, err := dynamic.NewForConfig(restConfig)
if err != nil {
return clientCacheEntry{}, fmt.Errorf("failed to create dynamic client: %w", err)
}
discoveryClient := memory.NewMemCacheClient(clientSet.Discovery())
mapper := restmapper.NewDeferredDiscoveryRESTMapper(discoveryClient)
cache := clientCacheEntry{
clientSet: clientSet,
dynamicClient: dynamicClient,
restConfig: restConfig,
discovery: discoveryClient,
mapper: mapper,
}
kubeInitMu.Lock()
defer kubeInitMu.Unlock()
if client, ok := clientCache[key]; ok {
return client, nil
if existingCache, ok := clientCache[key]; ok {
return existingCache, nil
}
initOpts := kube.InitOptions{
KubeConfigOptions: kube.KubeConfigOptions{
Context: kubeContext,
ConfigPath: kubeconfig,
},
}
clientCache[key] = cache
if err := kube.Init(initOpts); err != nil {
return nil, err
}
client := kube.Kubernetes
clientCache[key] = client
return client, nil
return cache, nil
}
func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Resource) error {
@@ -155,16 +234,29 @@ func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Reso
Namespace: namespace,
SkipLogs: !t.trackOptions.Logs,
})
case "canary":
specs.Canaries = append(specs.Canaries, multitrack.MultitrackSpec{
ResourceName: res.Name,
Namespace: namespace,
SkipLogs: !t.trackOptions.Logs,
})
default:
t.logger.Debugf("Skipping unsupported kind %s for resource %s/%s", res.Kind, namespace, res.Name)
}
}
if len(specs.Deployments)+len(specs.StatefulSets)+len(specs.DaemonSets)+len(specs.Jobs) == 0 {
t.logger.Info("No trackable resources found (only Deployment, StatefulSet, DaemonSet, and Job are supported)")
totalResources := len(specs.Deployments) + len(specs.StatefulSets) +
len(specs.DaemonSets) + len(specs.Jobs) + len(specs.Canaries)
if totalResources == 0 {
t.logger.Info("No trackable resources found (only Deployment, StatefulSet, DaemonSet, Job, and Canary are supported)")
return nil
}
t.logger.Infof("Tracking breakdown: Deployments=%d, StatefulSets=%d, DaemonSets=%d, Jobs=%d, Canaries=%d",
len(specs.Deployments), len(specs.StatefulSets), len(specs.DaemonSets),
len(specs.Jobs), len(specs.Canaries))
opts := multitrack.MultitrackOptions{
Options: tracker.Options{
ParentContext: ctx,
@@ -172,6 +264,9 @@ func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Reso
LogsFromTime: time.Now().Add(-t.trackOptions.LogsSince),
},
StatusProgressPeriod: 5 * time.Second,
DynamicClient: t.dynamicClient,
DiscoveryClient: t.discovery,
Mapper: t.mapper,
}
err := multitrack.Multitrack(t.clientSet, specs, opts)
+189
View File
@@ -1,10 +1,13 @@
package kubedog
import (
"math"
"os"
"testing"
"time"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/helmfile/helmfile/pkg/resource"
)
@@ -86,3 +89,189 @@ func TestTrackOptions_WithFilterConfig(t *testing.T) {
assert.Equal(t, []string{"Deployment", "StatefulSet"}, opts.Filter.TrackKinds)
assert.Equal(t, []string{"ConfigMap"}, opts.Filter.SkipKinds)
}
func TestTrackOptions_WithQPS(t *testing.T) {
opts := NewTrackOptions()
opts = opts.WithQPS(50.0)
assert.Equal(t, float32(50.0), opts.QPS)
}
func TestTrackOptions_WithBurst(t *testing.T) {
opts := NewTrackOptions()
opts = opts.WithBurst(100)
assert.Equal(t, 100, opts.Burst)
}
func TestTrackOptions_DefaultQPSBurst(t *testing.T) {
opts := NewTrackOptions()
assert.Equal(t, float32(100), opts.QPS)
assert.Equal(t, 200, opts.Burst)
}
func TestTrackerConfig_WithQPSBurst(t *testing.T) {
qps := float32(50.0)
burst := 100
config := &TrackerConfig{
Logger: nil,
Namespace: "test-ns",
KubeContext: "test-ctx",
Kubeconfig: "/test/kubeconfig",
TrackOptions: NewTrackOptions(),
KubedogQPS: &qps,
KubedogBurst: &burst,
}
assert.NotNil(t, config)
assert.Equal(t, "test-ns", config.Namespace)
assert.Equal(t, &qps, config.KubedogQPS)
assert.Equal(t, &burst, config.KubedogBurst)
assert.Equal(t, float32(50.0), *config.KubedogQPS)
assert.Equal(t, 100, *config.KubedogBurst)
}
func TestNewTracker_InvalidQPS(t *testing.T) {
invalidQPS := float32(-1.0)
burst := 100
cfg := &TrackerConfig{
Logger: nil,
Namespace: "test-ns",
KubeContext: "test-ctx",
Kubeconfig: "/nonexistent/kubeconfig",
TrackOptions: NewTrackOptions(),
KubedogQPS: &invalidQPS,
KubedogBurst: &burst,
}
tr, err := NewTracker(cfg)
assert.Error(t, err)
assert.Nil(t, tr)
assert.Contains(t, err.Error(), "invalid kubedog QPS")
assert.Contains(t, err.Error(), "must be > 0")
}
func TestNewTracker_NaNQPS(t *testing.T) {
nanQPS := float32(math.NaN())
burst := 100
cfg := &TrackerConfig{
Logger: nil,
Namespace: "test-ns",
KubeContext: "test-ctx",
Kubeconfig: "/nonexistent/kubeconfig",
TrackOptions: NewTrackOptions(),
KubedogQPS: &nanQPS,
KubedogBurst: &burst,
}
tr, err := NewTracker(cfg)
assert.Error(t, err)
assert.Nil(t, tr)
assert.Contains(t, err.Error(), "invalid kubedog QPS")
assert.Contains(t, err.Error(), "must be > 0 and finite")
}
func TestNewTracker_InfQPS(t *testing.T) {
infQPS := float32(math.Inf(1))
burst := 100
cfg := &TrackerConfig{
Logger: nil,
Namespace: "test-ns",
KubeContext: "test-ctx",
Kubeconfig: "/nonexistent/kubeconfig",
TrackOptions: NewTrackOptions(),
KubedogQPS: &infQPS,
KubedogBurst: &burst,
}
tr, err := NewTracker(cfg)
assert.Error(t, err)
assert.Nil(t, tr)
assert.Contains(t, err.Error(), "invalid kubedog QPS")
assert.Contains(t, err.Error(), "must be > 0 and finite")
}
func TestNewTracker_InvalidBurst(t *testing.T) {
qps := float32(50.0)
invalidBurst := 0
cfg := &TrackerConfig{
Logger: nil,
Namespace: "test-ns",
KubeContext: "test-ctx",
Kubeconfig: "/nonexistent/kubeconfig",
TrackOptions: NewTrackOptions(),
KubedogQPS: &qps,
KubedogBurst: &invalidBurst,
}
tr, err := NewTracker(cfg)
assert.Error(t, err)
assert.Nil(t, tr)
assert.Contains(t, err.Error(), "invalid kubedog burst")
assert.Contains(t, err.Error(), "must be >= 1")
}
func TestNewTracker_ValidQPSBurst(t *testing.T) {
qps := float32(50.0)
burst := 100
// Create a minimal valid kubeconfig in a temp file
tmpFile, err := os.CreateTemp("", "kubeconfig-*.yaml")
require.NoError(t, err)
defer os.Remove(tmpFile.Name())
kubeconfigContent := `
apiVersion: v1
kind: Config
clusters:
- cluster:
server: https://test-server:6443
name: test-cluster
contexts:
- context:
cluster: test-cluster
user: test-user
name: test-context
current-context: test-context
users:
- name: test-user
user:
token: test-token
`
_, err = tmpFile.WriteString(kubeconfigContent)
require.NoError(t, err)
require.NoError(t, tmpFile.Close())
cfg := &TrackerConfig{
Logger: nil,
Namespace: "test-ns",
KubeContext: "test-context",
Kubeconfig: tmpFile.Name(),
TrackOptions: NewTrackOptions(),
KubedogQPS: &qps,
KubedogBurst: &burst,
}
// This should succeed - validation passes and client is created
tr, err := NewTracker(cfg)
// The test should pass validation. It may fail later due to invalid cluster,
// but that's okay - we're testing that QPS/Burst validation works.
if err != nil {
// If there's an error, it should NOT be about invalid QPS/Burst
assert.NotContains(t, err.Error(), "invalid kubedog QPS")
assert.NotContains(t, err.Error(), "invalid kubedog burst")
} else {
// If no error, tracker should be created successfully
assert.NotNil(t, tr)
}
}
+2
View File
@@ -480,6 +480,8 @@ func (st *HelmState) trackWithKubedog(ctx context.Context, release *ReleaseSpec,
KubeContext: kubeContext,
Kubeconfig: st.kubeconfig,
TrackOptions: trackOpts,
KubedogQPS: release.KubedogQPS,
KubedogBurst: release.KubedogBurst,
})
if err != nil {
return fmt.Errorf("failed to create kubedog tracker: %w", err)
+4
View File
@@ -466,6 +466,10 @@ type ReleaseSpec struct {
SkipKinds []string `yaml:"skipKinds,omitempty"`
// TrackResources is a whitelist of specific resources to track
TrackResources []TrackResourceSpec `yaml:"trackResources,omitempty"`
// KubedogQPS specifies the QPS (queries per second) for kubedog kubernetes client
KubedogQPS *float32 `yaml:"kubedogQPS,omitempty"`
// KubedogBurst specifies the burst for kubedog kubernetes client
KubedogBurst *int `yaml:"kubedogBurst,omitempty"`
}
// TrackResourceSpec specifies a resource to track
+6 -6
View File
@@ -38,39 +38,39 @@ func TestGenerateID(t *testing.T) {
run(testcase{
subject: "baseline",
release: ReleaseSpec{Name: "foo", Chart: "incubator/raw"},
want: "foo-values-dd88b94b8",
want: "foo-values-6d799cf798",
})
run(testcase{
subject: "different bytes content",
release: ReleaseSpec{Name: "foo", Chart: "incubator/raw"},
data: []byte(`{"k":"v"}`),
want: "foo-values-6fb7bbb95f",
want: "foo-values-7f885447bf",
})
run(testcase{
subject: "different map content",
release: ReleaseSpec{Name: "foo", Chart: "incubator/raw"},
data: map[string]any{"k": "v"},
want: "foo-values-56d84c9897",
want: "foo-values-86f5d8fb55",
})
run(testcase{
subject: "different chart",
release: ReleaseSpec{Name: "foo", Chart: "stable/envoy"},
want: "foo-values-6644fc9d47",
want: "foo-values-5cd5c65db5",
})
run(testcase{
subject: "different name",
release: ReleaseSpec{Name: "bar", Chart: "incubator/raw"},
want: "bar-values-859cd849bf",
want: "bar-values-c59b4f979",
})
run(testcase{
subject: "specific ns",
release: ReleaseSpec{Name: "foo", Chart: "incubator/raw", Namespace: "myns"},
want: "myns-foo-values-86d544f7f9",
want: "myns-foo-values-56d6cd88cc",
})
for id, n := range ids {