diff --git a/cmd/controller/leader_election_test.go b/cmd/controller/leader_election_test.go new file mode 100644 index 0000000..35a214c --- /dev/null +++ b/cmd/controller/leader_election_test.go @@ -0,0 +1,129 @@ +// SPDX-License-Identifier: AGPL-3.0-only +package controller + +import ( + "context" + "net/http" + "net/http/httptest" + "sync" + "testing" + "time" + + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/runtime/serializer" + "k8s.io/client-go/rest" + "k8s.io/client-go/util/flowcontrol" +) + +type remoteAddrRecorder struct { + mu sync.Mutex + addrs []string +} + +func (r *remoteAddrRecorder) ServeHTTP(w http.ResponseWriter, req *http.Request) { + r.mu.Lock() + r.addrs = append(r.addrs, req.RemoteAddr) + r.mu.Unlock() + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{}`)) +} + +func (r *remoteAddrRecorder) last() string { + r.mu.Lock() + defer r.mu.Unlock() + return r.addrs[len(r.addrs)-1] +} + +func restClientFor(t *testing.T, cfg *rest.Config) *rest.RESTClient { + t.Helper() + cfg = rest.CopyConfig(cfg) + cfg.GroupVersion = &schema.GroupVersion{Group: "coordination.k8s.io", Version: "v1"} + cfg.APIPath = "/apis" + cfg.NegotiatedSerializer = serializer.NewCodecFactory(scheme).WithoutConversion() + c, err := rest.RESTClientFor(cfg) + if err != nil { + t.Fatalf("building REST client: %v", err) + } + return c +} + +func getLease(ctx context.Context, c *rest.RESTClient) error { + return c.Get().Namespace("default").Resource("leases").Name("81afa9db.datumapis.com").Do(ctx).Error() +} + +func TestLeaderElectionRestConfigDoesNotMutateBase(t *testing.T) { + limiter := flowcontrol.NewTokenBucketRateLimiter(1, 1) + base := &rest.Config{Host: "https://example.invalid", QPS: 50, Burst: 100, RateLimiter: limiter} + + cfg := leaderElectionRestConfig(base) + + if cfg.RateLimiter != nil { + t.Errorf("leader election config reuses the base rate limiter") + } + if cfg.QPS != leaderElectionQPS || cfg.Burst != leaderElectionBurst { + t.Errorf("got QPS %v Burst %d, want %v and %d", cfg.QPS, cfg.Burst, leaderElectionQPS, leaderElectionBurst) + } + if cfg.Dial == nil { + t.Errorf("leader election config has no dialer of its own") + } + if base.RateLimiter != limiter || base.QPS != 50 || base.Burst != 100 || base.Dial != nil { + t.Errorf("base config was mutated: %+v", base) + } +} + +func TestLeaderElectionIgnoresSaturatedMainLimiter(t *testing.T) { + srv := httptest.NewServer(&remoteAddrRecorder{}) + defer srv.Close() + + base := &rest.Config{ + Host: srv.URL, + RateLimiter: flowcontrol.NewTokenBucketRateLimiter(0.001, 1), + } + mainClient := restClientFor(t, base) + + if err := getLease(context.Background(), mainClient); err != nil { + t.Fatalf("first request on main client: %v", err) + } + + ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond) + defer cancel() + if err := getLease(ctx, mainClient); err == nil { + t.Fatalf("expected the main client to be throttled") + } + + leaderClient := restClientFor(t, leaderElectionRestConfig(base)) + ctx, cancel = context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + if err := getLease(ctx, leaderClient); err != nil { + t.Fatalf("leader election request was throttled by the main limiter: %v", err) + } +} + +func TestLeaderElectionUsesItsOwnConnection(t *testing.T) { + rec := &remoteAddrRecorder{} + srv := httptest.NewServer(rec) + defer srv.Close() + + base := &rest.Config{Host: srv.URL, QPS: -1} + mainClient := restClientFor(t, base) + ctx := context.Background() + + if err := getLease(ctx, mainClient); err != nil { + t.Fatalf("main client request: %v", err) + } + mainAddr := rec.last() + + if err := getLease(ctx, restClientFor(t, base)); err != nil { + t.Fatalf("second main client request: %v", err) + } + if rec.last() != mainAddr { + t.Fatalf("clients built from the same config did not share a connection, so this test cannot tell connections apart") + } + + if err := getLease(ctx, restClientFor(t, leaderElectionRestConfig(base))); err != nil { + t.Fatalf("leader election request: %v", err) + } + if rec.last() == mainAddr { + t.Errorf("leader election reused the main client's connection %s", mainAddr) + } +} diff --git a/cmd/controller/manager.go b/cmd/controller/manager.go index 9b3c6af..6419aa3 100644 --- a/cmd/controller/manager.go +++ b/cmd/controller/manager.go @@ -5,6 +5,7 @@ import ( "crypto/tls" "flag" "fmt" + "net" "os" "path/filepath" "time" @@ -18,6 +19,7 @@ import ( utilruntime "k8s.io/apimachinery/pkg/util/runtime" utilfeature "k8s.io/apiserver/pkg/util/feature" clientgoscheme "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/rest" cliflag "k8s.io/component-base/cli/flag" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/certwatcher" @@ -114,7 +116,7 @@ func NewControllerManagerCommand() *cobra.Command { "The duration that the acting leader will retry refreshing leadership before giving up.") cmd.Flags().DurationVar(&leaderElectionRetryPeriod, "leader-election-retry-period", 2*time.Second, "The duration the LeaderElector clients should wait between tries of actions.") - cmd.Flags().BoolVar(&leaderElectionReleaseOnCancel, "leader-election-release-on-cancel", false, + cmd.Flags().BoolVar(&leaderElectionReleaseOnCancel, "leader-election-release-on-cancel", true, "If the leader should step down voluntarily when the Manager ends. "+ "This requires the binary to immediately end when the Manager is stopped.") @@ -280,7 +282,8 @@ func runControllerManager( }) } - mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{ + restConfig := ctrl.GetConfigOrDie() + mgr, err := ctrl.NewManager(restConfig, ctrl.Options{ Scheme: scheme, Metrics: metricsServerOptions, WebhookServer: webhookServer, @@ -288,6 +291,7 @@ func runControllerManager( LeaderElection: enableLeaderElection, LeaderElectionID: leaderElectionID, LeaderElectionNamespace: leaderElectionNamespace, + LeaderElectionConfig: leaderElectionRestConfig(restConfig), LeaseDuration: &leaderElectionLeaseDuration, RenewDeadline: &leaderElectionRenewDeadline, RetryPeriod: &leaderElectionRetryPeriod, @@ -347,3 +351,20 @@ func runControllerManager( return nil } + +const ( + leaderElectionQPS = 5 + leaderElectionBurst = 10 +) + +func leaderElectionRestConfig(base *rest.Config) *rest.Config { + cfg := rest.CopyConfig(base) + cfg.RateLimiter = nil + cfg.QPS = leaderElectionQPS + cfg.Burst = leaderElectionBurst + cfg.Dial = (&net.Dialer{ + Timeout: 30 * time.Second, + KeepAlive: 30 * time.Second, + }).DialContext + return cfg +} diff --git a/config/manager/manager.yaml b/config/manager/manager.yaml index 3c7eb4a..1db65f2 100644 --- a/config/manager/manager.yaml +++ b/config/manager/manager.yaml @@ -86,7 +86,7 @@ spec: - name: LEADER_ELECTION_RETRY_PERIOD value: "2s" - name: LEADER_ELECTION_RELEASE_ON_CANCEL - value: "false" + value: "true" - name: METRICS_SECURE value: "true" - name: WEBHOOK_CERT_PATH