This is an automated email from the ASF dual-hosted git repository.
AlinsRan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/apisix-ingress-controller.git
The following commit(s) were added to refs/heads/master by this push:
new a1c2ec17 fix: never treat an unreachable API server as a missing API
resource (#2817)
a1c2ec17 is described below
commit a1c2ec172b75a10d5754b05c9b4700cd5360ca87
Author: AlinsRan <[email protected]>
AuthorDate: Tue Jul 28 14:48:30 2026 +0800
fix: never treat an unreachable API server as a missing API resource (#2817)
---
internal/controller/consumer_controller.go | 9 +-
internal/controller/gateway_controller.go | 18 +-
internal/controller/gatewayproxy_controller.go | 9 +-
internal/controller/indexer/indexer.go | 42 ++--
internal/controller/tcproute_controller.go | 6 +-
internal/controller/tlsroute_controller.go | 6 +-
internal/controller/udproute_controller.go | 6 +-
internal/manager/controllers.go | 100 +++++---
internal/manager/run.go | 71 +++++-
pkg/utils/k8s.go | 204 ++++++++++++++--
pkg/utils/k8s_test.go | 324 +++++++++++++++++++++++++
11 files changed, 705 insertions(+), 90 deletions(-)
diff --git a/internal/controller/consumer_controller.go
b/internal/controller/consumer_controller.go
index 2c0a4eb0..0f265be5 100644
--- a/internal/controller/consumer_controller.go
+++ b/internal/controller/consumer_controller.go
@@ -62,7 +62,14 @@ type ConsumerReconciler struct { //nolint:revive
// SetupWithManager sets up the controller with the Manager.
func (r *ConsumerReconciler) SetupWithManager(mgr ctrl.Manager) error {
- if config.ControllerConfig.DisableGatewayAPI ||
!pkgutils.HasAPIResource(mgr, &gatewayv1.Gateway{}) {
+ hasGatewayAPI := false
+ if !config.ControllerConfig.DisableGatewayAPI {
+ var err error
+ if hasGatewayAPI, err = pkgutils.HasAPIResource(mgr,
&gatewayv1.Gateway{}); err != nil {
+ return err
+ }
+ }
+ if !hasGatewayAPI {
r.Log.Info("skipping Consumer controller setup as Gateway API
is not available")
return nil
}
diff --git a/internal/controller/gateway_controller.go
b/internal/controller/gateway_controller.go
index 5b5130d7..69cc89f5 100644
--- a/internal/controller/gateway_controller.go
+++ b/internal/controller/gateway_controller.go
@@ -105,19 +105,31 @@ func (r *GatewayReconciler) SetupWithManager(mgr
ctrl.Manager) error {
builder.WithPredicates(referenceGrantPredicates(KindGateway)),
)
}
- if pkgutils.HasAPIResource(mgr, &gatewayv1alpha2.TCPRoute{}) {
+ hasTCPRoute, err := pkgutils.HasAPIResource(mgr,
&gatewayv1alpha2.TCPRoute{})
+ if err != nil {
+ return err
+ }
+ if hasTCPRoute {
bdr.Watches(
&gatewayv1alpha2.TCPRoute{},
handler.EnqueueRequestsFromMapFunc(r.listGatewaysForStatusParentRefs),
)
}
- if pkgutils.HasAPIResource(mgr, &gatewayv1alpha2.TLSRoute{}) {
+ hasTLSRoute, err := pkgutils.HasAPIResource(mgr,
&gatewayv1alpha2.TLSRoute{})
+ if err != nil {
+ return err
+ }
+ if hasTLSRoute {
bdr.Watches(
&gatewayv1alpha2.TLSRoute{},
handler.EnqueueRequestsFromMapFunc(r.listGatewaysForStatusParentRefs),
)
}
- if pkgutils.HasAPIResource(mgr, &gatewayv1alpha2.UDPRoute{}) {
+ hasUDPRoute, err := pkgutils.HasAPIResource(mgr,
&gatewayv1alpha2.UDPRoute{})
+ if err != nil {
+ return err
+ }
+ if hasUDPRoute {
bdr.Watches(
&gatewayv1alpha2.UDPRoute{},
handler.EnqueueRequestsFromMapFunc(r.listGatewaysForStatusParentRefs),
diff --git a/internal/controller/gatewayproxy_controller.go
b/internal/controller/gatewayproxy_controller.go
index 82bb2c71..ab53ad6f 100644
--- a/internal/controller/gatewayproxy_controller.go
+++ b/internal/controller/gatewayproxy_controller.go
@@ -54,9 +54,14 @@ type GatewayProxyController struct {
}
func (r *GatewayProxyController) SetupWithManager(mrg ctrl.Manager) error {
- if config.ControllerConfig.DisableGatewayAPI ||
!pkgutils.HasAPIResource(mrg, &gatewayv1.Gateway{}) {
- r.disableGatewayAPI = true
+ hasGatewayAPI := false
+ if !config.ControllerConfig.DisableGatewayAPI {
+ var err error
+ if hasGatewayAPI, err = pkgutils.HasAPIResource(mrg,
&gatewayv1.Gateway{}); err != nil {
+ return err
+ }
}
+ r.disableGatewayAPI = !hasGatewayAPI
builder := ctrl.NewControllerManagedBy(mrg).
For(&v1alpha1.GatewayProxy{}).
WithEventFilter(
diff --git a/internal/controller/indexer/indexer.go
b/internal/controller/indexer/indexer.go
index 037e9e11..99f8b535 100644
--- a/internal/controller/indexer/indexer.go
+++ b/internal/controller/indexer/indexer.go
@@ -65,12 +65,16 @@ func SetupAPIv1alpha1Indexer(mgr ctrl.Manager) error {
&v1alpha1.GatewayProxy{}: setupGatewayProxyIndexer,
&v1alpha1.L4RoutePolicy{}: setupL4RoutePolicyIndexer,
} {
- if utils.HasAPIResource(mgr, resource) {
- if err := setup(mgr); err != nil {
- return err
- }
- } else {
+ installed, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return err
+ }
+ if !installed {
setupLog.Info("Skipping indexer setup, API not found in
cluster", "api", utils.FormatGVK(resource))
+ continue
+ }
+ if err := setup(mgr); err != nil {
+ return err
}
}
return nil
@@ -87,12 +91,16 @@ func SetupAPIv2Indexer(mgr ctrl.Manager) error {
&apiv2.ApisixTls{}: setupApisixTlsIndexer,
&apiv2.ApisixGlobalRule{}: setupApisixGlobalRuleIndexer,
} {
- if utils.HasAPIResource(mgr, resource) {
- if err := setup(mgr); err != nil {
- return err
- }
- } else {
+ installed, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return err
+ }
+ if !installed {
setupLog.Info("Skipping indexer setup, API not found in
cluster", "api", utils.FormatGVK(resource))
+ continue
+ }
+ if err := setup(mgr); err != nil {
+ return err
}
}
return nil
@@ -110,12 +118,16 @@ func SetupGatewayAPIIndexer(mgr ctrl.Manager) error {
&gatewayv1alpha2.TLSRoute{}: setupTLSRouteIndexer,
&gatewayv1.GatewayClass{}: setupGatewayClassIndexer,
} {
- if utils.HasAPIResource(mgr, resource) {
- if err := setup(mgr); err != nil {
- return err
- }
- } else {
+ installed, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return err
+ }
+ if !installed {
setupLog.Info("Skipping indexer setup, API not found in
cluster", "api", utils.FormatGVK(resource))
+ continue
+ }
+ if err := setup(mgr); err != nil {
+ return err
}
}
return nil
diff --git a/internal/controller/tcproute_controller.go
b/internal/controller/tcproute_controller.go
index f76f6855..a2360157 100644
--- a/internal/controller/tcproute_controller.go
+++ b/internal/controller/tcproute_controller.go
@@ -101,7 +101,11 @@ func (r *TCPRouteReconciler) SetupWithManager(mgr
ctrl.Manager) error {
// L4RoutePolicy is an optional CRD. Only watch it when installed so the
// controller still starts if the CRD has not been applied yet (e.g.
upgrades).
- r.supportsL4RoutePolicy = pkgutils.HasAPIResource(mgr,
&v1alpha1.L4RoutePolicy{})
+ supportsL4RoutePolicy, err := pkgutils.HasAPIResource(mgr,
&v1alpha1.L4RoutePolicy{})
+ if err != nil {
+ return err
+ }
+ r.supportsL4RoutePolicy = supportsL4RoutePolicy
if r.supportsL4RoutePolicy {
bdr.Watches(&v1alpha1.L4RoutePolicy{},
handler.EnqueueRequestsFromMapFunc(r.listTCPRoutesForL4RoutePolicy),
diff --git a/internal/controller/tlsroute_controller.go
b/internal/controller/tlsroute_controller.go
index dd6b7ff4..edb7381f 100644
--- a/internal/controller/tlsroute_controller.go
+++ b/internal/controller/tlsroute_controller.go
@@ -101,7 +101,11 @@ func (r *TLSRouteReconciler) SetupWithManager(mgr
ctrl.Manager) error {
// L4RoutePolicy is an optional CRD. Only watch it when installed so the
// controller still starts if the CRD has not been applied yet (e.g.
upgrades).
- r.supportsL4RoutePolicy = pkgutils.HasAPIResource(mgr,
&v1alpha1.L4RoutePolicy{})
+ supportsL4RoutePolicy, err := pkgutils.HasAPIResource(mgr,
&v1alpha1.L4RoutePolicy{})
+ if err != nil {
+ return err
+ }
+ r.supportsL4RoutePolicy = supportsL4RoutePolicy
if r.supportsL4RoutePolicy {
bdr.Watches(&v1alpha1.L4RoutePolicy{},
handler.EnqueueRequestsFromMapFunc(r.listTLSRoutesForL4RoutePolicy),
diff --git a/internal/controller/udproute_controller.go
b/internal/controller/udproute_controller.go
index 33c698ae..17fb5f6a 100644
--- a/internal/controller/udproute_controller.go
+++ b/internal/controller/udproute_controller.go
@@ -101,7 +101,11 @@ func (r *UDPRouteReconciler) SetupWithManager(mgr
ctrl.Manager) error {
// L4RoutePolicy is an optional CRD. Only watch it when installed so the
// controller still starts if the CRD has not been applied yet (e.g.
upgrades).
- r.supportsL4RoutePolicy = pkgutils.HasAPIResource(mgr,
&v1alpha1.L4RoutePolicy{})
+ supportsL4RoutePolicy, err := pkgutils.HasAPIResource(mgr,
&v1alpha1.L4RoutePolicy{})
+ if err != nil {
+ return err
+ }
+ r.supportsL4RoutePolicy = supportsL4RoutePolicy
if r.supportsL4RoutePolicy {
bdr.Watches(&v1alpha1.L4RoutePolicy{},
handler.EnqueueRequestsFromMapFunc(r.listUDPRoutesForL4RoutePolicy),
diff --git a/internal/manager/controllers.go b/internal/manager/controllers.go
index 5e48b397..af518e3c 100644
--- a/internal/manager/controllers.go
+++ b/internal/manager/controllers.go
@@ -211,11 +211,15 @@ func setupGatewayAPIControllers(ctx context.Context, mgr
manager.Manager, pro pr
Readier: readier,
},
} {
- if utils.HasAPIResource(mgr, resource) {
- runnables = append(runnables, controller)
- } else {
- setupLog.Info("Skipping indexer setup, API not found in
cluster", "api", utils.FormatGVK(resource))
+ installed, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return nil, err
}
+ if !installed {
+ setupLog.Info("Skipping controller setup, API not found
in cluster", "api", utils.FormatGVK(resource))
+ continue
+ }
+ runnables = append(runnables, controller)
}
return runnables, nil
}
@@ -288,69 +292,88 @@ func setupAPIv2Controllers(ctx context.Context, mgr
manager.Manager, pro provide
Updater: updater,
},
} {
- if utils.HasAPIResource(mgr, resource) {
- runnables = append(runnables, controller)
- } else {
- setupLog.Info("Skipping indexer setup, API not found in
cluster", "api", utils.FormatGVK(resource))
+ installed, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return nil, err
+ }
+ if !installed {
+ setupLog.Info("Skipping controller setup, API not found
in cluster", "api", utils.FormatGVK(resource))
+ continue
}
+ runnables = append(runnables, controller)
}
return runnables, nil
}
-func registerReadiness(mgr manager.Manager, readier
readiness.ReadinessManager) {
+func registerReadiness(mgr manager.Manager, readier
readiness.ReadinessManager) error {
log := ctrl.LoggerFrom(context.Background()).WithName("readiness")
- registerAPIv2ForReadiness(mgr, log, readier)
+ if err := registerAPIv2ForReadiness(mgr, log, readier); err != nil {
+ return err
+ }
if !config.ControllerConfig.DisableGatewayAPI {
- registerGatewayAPIForReadiness(mgr, log, readier)
+ if err := registerGatewayAPIForReadiness(mgr, log, readier);
err != nil {
+ return err
+ }
}
- registerAPIv1alpha1ForReadiness(mgr, log, readier)
+ return registerAPIv1alpha1ForReadiness(mgr, log, readier)
}
func registerGatewayAPIForReadiness(
mgr manager.Manager,
log logr.Logger,
readier readiness.ReadinessManager,
-) {
- var installed []schema.GroupVersionKind
- for _, resource := range []client.Object{
+) error {
+ resources := []client.Object{
&gatewayv1.HTTPRoute{},
&gatewayv1.GRPCRoute{},
&gatewayv1alpha2.TCPRoute{},
&gatewayv1alpha2.UDPRoute{},
&gatewayv1alpha2.TLSRoute{},
- } {
+ }
+ installed := make([]schema.GroupVersionKind, 0, len(resources))
+ for _, resource := range resources {
gvk := types.GvkOf(resource)
- if utils.HasAPIResource(mgr, resource) {
- installed = append(installed, gvk)
- } else {
+ has, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return err
+ }
+ if !has {
log.Info("Skipping readiness registration, API not
found", "gvk", gvk)
+ continue
}
+ installed = append(installed, gvk)
}
if len(installed) == 0 {
- return
+ return nil
}
readier.RegisterGVK(readiness.GVKConfig{GVKs: installed})
+ return nil
}
func registerAPIv2ForReadiness(
mgr manager.Manager,
log logr.Logger,
readier readiness.ReadinessManager,
-) {
- var installed []schema.GroupVersionKind
- for _, resource := range apiV2ReadinessResources() {
+) error {
+ resources := apiV2ReadinessResources()
+ installed := make([]schema.GroupVersionKind, 0, len(resources))
+ for _, resource := range resources {
gvk := types.GvkOf(resource)
- if utils.HasAPIResource(mgr, resource) {
- installed = append(installed, gvk)
- } else {
+ has, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return err
+ }
+ if !has {
log.Info("Skipping readiness registration, API not
found", "gvk", gvk)
+ continue
}
+ installed = append(installed, gvk)
}
if len(installed) == 0 {
- return
+ return nil
}
readier.RegisterGVK(readiness.GVKConfig{
@@ -361,6 +384,7 @@ func registerAPIv2ForReadiness(
return ingressClass != nil
}),
})
+ return nil
}
func apiV2ReadinessResources() []client.Object {
@@ -377,20 +401,25 @@ func registerAPIv1alpha1ForReadiness(
mgr manager.Manager,
log logr.Logger,
readier readiness.ReadinessManager,
-) {
- var installed []schema.GroupVersionKind
- for _, resource := range []client.Object{
+) error {
+ resources := []client.Object{
&v1alpha1.Consumer{},
- } {
+ }
+ installed := make([]schema.GroupVersionKind, 0, len(resources))
+ for _, resource := range resources {
gvk := types.GvkOf(resource)
- if utils.HasAPIResource(mgr, resource) {
- installed = append(installed, gvk)
- } else {
+ has, err := utils.HasAPIResource(mgr, resource)
+ if err != nil {
+ return err
+ }
+ if !has {
log.Info("Skipping readiness registration, API not
found", "gvk", gvk)
+ continue
}
+ installed = append(installed, gvk)
}
if len(installed) == 0 {
- return
+ return nil
}
readier.RegisterGVK(readiness.GVKConfig{
@@ -403,4 +432,5 @@ func registerAPIv1alpha1ForReadiness(
return
controller.MatchConsumerGatewayRef(context.Background(), mgr.GetClient(), log,
consumer)
}),
})
+ return nil
}
diff --git a/internal/manager/run.go b/internal/manager/run.go
index 315644da..34c2d338 100644
--- a/internal/manager/run.go
+++ b/internal/manager/run.go
@@ -24,7 +24,6 @@ import (
"github.com/go-logr/logr"
"k8s.io/apimachinery/pkg/runtime"
- "k8s.io/apimachinery/pkg/runtime/schema"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
"k8s.io/apimachinery/pkg/util/version"
"k8s.io/client-go/discovery"
@@ -49,6 +48,7 @@ import (
"github.com/apache/apisix-ingress-controller/internal/provider"
_ "github.com/apache/apisix-ingress-controller/internal/provider/init"
_ "github.com/apache/apisix-ingress-controller/pkg/metrics"
+ "github.com/apache/apisix-ingress-controller/pkg/utils"
)
var (
@@ -86,6 +86,11 @@ func Run(ctx context.Context, logger logr.Logger) error {
setupLog := ctrl.LoggerFrom(ctx).WithName("setup")
+ // SetupSignalHandler must be called exactly once. Install it before
setup
+ // starts so that a shutdown signal is honored while waiting for the API
+ // server, not only once the manager is running.
+ signalCtx := ctrl.SetupSignalHandler()
+
// if the enable-http2 flag is false (the default), http/2 should be
disabled
// due to its vulnerabilities. More specifically, disabling http/2 will
// prevent from being vulnerable to the HTTP/2 Stream Cancellation and
@@ -171,11 +176,40 @@ func Run(ctx context.Context, logger logr.Logger) error {
return err
}
+ // Which API resources are installed is detected once, below, and gates
the
+ // registration of field indexes, controllers and readiness checks for
the
+ // whole lifetime of the process. Detecting them against an unreachable
API
+ // server would silently classify everything as "not installed" and
leave the
+ // controller permanently degraded until it is restarted, so wait for
the API
+ // server to answer first.
+ if err := utils.WaitForAPIServer(signalCtx, mgr.GetConfig(), setupLog);
err != nil {
+ setupLog.Error(err, "unable to reach the Kubernetes API server")
+ return err
+ }
+
+ // API resource detection runs outside the manager and does not observe
+ // signalCtx, so check between setup phases: now that the signal
handler is
+ // installed this early, a shutdown signal would otherwise be swallowed
until
+ // mgr.Start.
+ checkShutdown := func() error {
+ if err := signalCtx.Err(); err != nil {
+ setupLog.Info("shutdown requested during setup,
stopping", "reason", err)
+ return err
+ }
+ return nil
+ }
+
// Check Kubernetes cluster version
checkK8sVersion(mgr, setupLog)
readier := readiness.NewReadinessManager(mgr.GetClient(), logger)
- registerReadiness(mgr, readier)
+ if err := registerReadiness(mgr, readier); err != nil {
+ setupLog.Error(err, "unable to register readiness checks")
+ return err
+ }
+ if err := checkShutdown(); err != nil {
+ return err
+ }
if err := mgr.Add(readier); err != nil {
setupLog.Error(err, "unable to add readiness manager")
@@ -214,16 +248,26 @@ func Run(ctx context.Context, logger logr.Logger) error {
return err
}
- setupLog.Info("check ReferenceGrants is enabled")
- _, err = mgr.GetRESTMapper().KindsFor(schema.GroupVersionResource{
- Group: v1beta1.GroupVersion.Group,
- Version: v1beta1.GroupVersion.Version,
- Resource: "referencegrants",
- })
- if err != nil {
- setupLog.Info("CRD ReferenceGrants is not installed", "err",
err)
+ // ReferenceGrant is a Gateway API kind and only consulted by Gateway
API
+ // paths, so skip the detection entirely when Gateway API is disabled.
+ hasReferenceGrant := false
+ if config.ControllerConfig.DisableGatewayAPI {
+ setupLog.Info("Gateway API is disabled, skipping the
ReferenceGrants check")
+ } else {
+ setupLog.Info("check ReferenceGrants is enabled")
+ if hasReferenceGrant, err = utils.HasAPIResource(mgr,
&v1beta1.ReferenceGrant{}); err != nil {
+ setupLog.Error(err, "unable to detect whether
ReferenceGrants is installed")
+ return err
+ }
+ if !hasReferenceGrant {
+ setupLog.Info("CRD ReferenceGrants is not installed,
cross-namespace references will be rejected",
+ "gvk",
utils.FormatGVK(&v1beta1.ReferenceGrant{}))
+ }
+ }
+ controller.SetEnableReferenceGrant(hasReferenceGrant)
+ if err := checkShutdown(); err != nil {
+ return err
}
- controller.SetEnableReferenceGrant(err == nil)
setupLog.Info("setting up controllers")
controllers, err := setupControllers(ctx, mgr, provider,
updater.Writer(), readier)
@@ -236,6 +280,9 @@ func Run(ctx context.Context, logger logr.Logger) error {
if err := c.SetupWithManager(mgr); err != nil {
return err
}
+ if err := checkShutdown(); err != nil {
+ return err
+ }
}
// +kubebuilder:scaffold:builder
@@ -263,7 +310,7 @@ func Run(ctx context.Context, logger logr.Logger) error {
}
setupLog.Info("starting controller manager")
- return mgr.Start(ctrl.SetupSignalHandler())
+ return mgr.Start(signalCtx)
}
func checkK8sVersion(mgr ctrl.Manager, logger logr.Logger) {
diff --git a/pkg/utils/k8s.go b/pkg/utils/k8s.go
index 425f2ff9..97c46e2b 100644
--- a/pkg/utils/k8s.go
+++ b/pkg/utils/k8s.go
@@ -18,8 +18,18 @@
package utils
import (
+ "context"
+ "errors"
+ "fmt"
+ "time"
+
"github.com/go-logr/logr"
+ apierrors "k8s.io/apimachinery/pkg/api/errors"
+ metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+ "k8s.io/apimachinery/pkg/runtime/schema"
+ "k8s.io/apimachinery/pkg/util/wait"
"k8s.io/client-go/discovery"
+ "k8s.io/client-go/rest"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/apiutil"
@@ -27,22 +37,81 @@ import (
"github.com/apache/apisix-ingress-controller/internal/types"
)
-// HasAPIResource checks if a specific API resource is available in the
current cluster.
-// It uses the Discovery API to query the cluster's available resources and
returns true
-// if the resource is found, false otherwise.
-func HasAPIResource(mgr ctrl.Manager, obj client.Object) bool {
+const (
+ // A discovery request that gives no definitive answer is retried a few
times
+ // to absorb short blips. A longer outage is reported to the caller
instead:
+ // it is not this function's job to decide how long to wait for the API
server.
+ discoveryRetryInitialInterval = 500 * time.Millisecond
+ discoveryRetrySteps = 5 // ~7.5s, up to ~8.3s with jitter
+
+ // The startup wait for the API server is deliberately finite. The
health
+ // probe endpoint is only served once mgr.Start runs, so a pod still
waiting
+ // here fails its liveness probe and is killed anyway; giving up with an
+ // explicit error and letting the pod restart says more in the logs than
+ // being killed mid-wait. Longer outages are absorbed by the restart
backoff.
+ apiServerWaitInitialInterval = 1 * time.Second
+ apiServerWaitSteps = 6 // ~31s, up to ~34s with jitter
+
+ backoffFactor = 2
+ backoffJitter = 0.1
+)
+
+// Neither backoff sets Cap: wait.Backoff zeroes the remaining steps as soon as
+// the capped interval is reached, which would silently cut the retries short.
+
+func discoveryBackoff() wait.Backoff {
+ return wait.Backoff{
+ Duration: discoveryRetryInitialInterval,
+ Factor: backoffFactor,
+ Jitter: backoffJitter,
+ Steps: discoveryRetrySteps,
+ }
+}
+
+func apiServerWaitBackoff() wait.Backoff {
+ return wait.Backoff{
+ Duration: apiServerWaitInitialInterval,
+ Factor: backoffFactor,
+ Jitter: backoffJitter,
+ Steps: apiServerWaitSteps,
+ }
+}
+
+// HasAPIResource reports whether the API resource of obj is served by the
cluster.
+//
+// Only a 404 from discovery reports the resource as absent. Every other
failure
+// returns an error, including a Forbidden discovery request: callers must not
+// fall back to "resource absent" there. Detection runs once at startup and
gates
+// field index, controller and readiness registration, so an API server outage
or
+// an RBAC mistake would otherwise leave the controller permanently degraded
+// until it is restarted.
+func HasAPIResource(mgr ctrl.Manager, obj client.Object) (bool, error) {
return HasAPIResourceWithLogger(mgr, obj,
ctrl.Log.WithName("api-detection"))
}
// HasAPIResourceWithLogger is the same as HasAPIResource but accepts a custom
logger
// for more detailed debugging information.
-func HasAPIResourceWithLogger(mgr ctrl.Manager, obj client.Object, logger
logr.Logger) bool {
+func HasAPIResourceWithLogger(mgr ctrl.Manager, obj client.Object, logger
logr.Logger) (bool, error) {
gvk, err := apiutil.GVKForObject(obj, mgr.GetScheme())
if err != nil {
- logger.Info("cannot derive GVK from scheme", "error", err)
- return false
+ return false, fmt.Errorf("cannot derive GVK from scheme: %w",
err)
+ }
+
+ // Create discovery client
+ discoveryClient, err :=
discovery.NewDiscoveryClientForConfig(mgr.GetConfig())
+ if err != nil {
+ return false, fmt.Errorf("failed to create discovery client:
%w", err)
}
+ return hasAPIResource(discoveryClient, gvk, discoveryBackoff(), logger)
+}
+
+func hasAPIResource(
+ discoveryClient discovery.DiscoveryInterface,
+ gvk schema.GroupVersionKind,
+ backoff wait.Backoff,
+ logger logr.Logger,
+) (bool, error) {
groupVersion := gvk.GroupVersion().String()
logger = logger.WithValues(
@@ -52,29 +121,126 @@ func HasAPIResourceWithLogger(mgr ctrl.Manager, obj
client.Object, logger logr.L
"groupVersion", groupVersion,
)
- // Create discovery client
- discoveryClient, err :=
discovery.NewDiscoveryClientForConfig(mgr.GetConfig())
- if err != nil {
- logger.Info("failed to create discovery client", "error", err)
- return false
- }
-
// Query server resources for the specific group/version
- apiResources, err :=
discoveryClient.ServerResourcesForGroupVersion(groupVersion)
- if err != nil {
+ var apiResources *metav1.APIResourceList
+ err := retryUntilDefinitive(backoff, logger, func() error {
+ var err error
+ apiResources, err =
discoveryClient.ServerResourcesForGroupVersion(groupVersion)
+ return err
+ })
+ switch {
+ case err == nil:
+ case apierrors.IsNotFound(err):
logger.Info("group/version not available in cluster", "error",
err)
- return false
+ return false, nil
+ case apierrors.IsForbidden(err):
+ // Definitive, but no evidence of absence: the resource may
well be served
+ // and only discovery is denied. Resolving it to "absent" would
be cached
+ // for the lifetime of the manager and permanently skip the
controller, so
+ // an RBAC misconfiguration must abort startup instead.
+ return false, fmt.Errorf("discovery of %s is forbidden, check
the discovery "+
+ "permissions of the controller service account: %w",
gvk, err)
+ default:
+ return false, fmt.Errorf("failed to detect API resource %s:
%w", gvk, err)
}
// Check if the specific kind exists in the resource list
for _, res := range apiResources.APIResources {
if res.Kind == gvk.Kind {
- return true
+ return true, nil
}
}
logger.Info("API resource kind not found in group/version")
- return false
+ return false, nil
+}
+
+// WaitForAPIServer blocks until the API server answers a discovery request,
ctx
+// is done, or the wait budget is exhausted.
+//
+// Capability detection (field indexes, optional CRDs, ReferenceGrant support)
is
+// decided once at startup and cannot be revised afterwards, so it must not run
+// against an unreachable API server.
+func WaitForAPIServer(ctx context.Context, cfg *rest.Config, logger
logr.Logger) error {
+ discoveryClient, err := discovery.NewDiscoveryClientForConfig(cfg)
+ if err != nil {
+ return fmt.Errorf("failed to create discovery client: %w", err)
+ }
+ return waitForAPIServer(ctx, func() error {
+ _, err := discoveryClient.ServerVersion()
+ return err
+ }, apiServerWaitBackoff(), logger)
+}
+
+func waitForAPIServer(ctx context.Context, probe func() error, backoff
wait.Backoff, logger logr.Logger) error {
+ var lastErr error
+ if err := wait.ExponentialBackoffWithContext(ctx, backoff,
func(context.Context) (bool, error) {
+ if lastErr = probe(); !apiServerAnswered(lastErr) {
+ logger.Info("waiting for the Kubernetes API server to
become reachable", "error", lastErr)
+ return false, nil
+ }
+ if lastErr != nil {
+ logger.Info("the Kubernetes API server rejected the
probe but is reachable, continuing", "error", lastErr)
+ }
+ return true, nil
+ }); err != nil {
+ // ExponentialBackoffWithContext only ever returns ctx.Err() or
its own
+ // interrupted error, so this tells cancellation and an
exhausted budget
+ // apart without racing against ctx.
+ if !errors.Is(err, context.Canceled) && !errors.Is(err,
context.DeadlineExceeded) && lastErr != nil {
+ err = lastErr
+ }
+ return fmt.Errorf("give up waiting for the Kubernetes API
server: %w", err)
+ }
+ return nil
+}
+
+// apiServerAnswered reports whether err still proves the API server responded.
+// An authn/authz rejection is an answer: a hardened cluster may not grant the
+// controller's service account access to the probed endpoint, and that is no
+// reason to hold up startup.
+func apiServerAnswered(err error) bool {
+ return err == nil || apierrors.IsUnauthorized(err) ||
apierrors.IsForbidden(err)
+}
+
+// retryUntilDefinitive runs probe until it returns nil or a definitive answer,
+// and reports the last failure once the retry budget is exhausted.
+func retryUntilDefinitive(backoff wait.Backoff, logger logr.Logger, probe
func() error) error {
+ var probed bool
+ var lastErr error
+ err := wait.ExponentialBackoff(backoff, func() (bool, error) {
+ probed = true
+ if lastErr = probe(); lastErr == nil || isDefinitive(lastErr) {
+ return true, nil
+ }
+ logger.Info("discovery request failed, retrying", "error",
lastErr)
+ return false, nil
+ })
+ if !probed {
+ // A backoff with no steps left never runs the condition.
Reporting nil
+ // here would be read as "the resource is present" by every
caller.
+ return fmt.Errorf("discovery was never attempted: %w", err)
+ }
+ return lastErr
+}
+
+// isDefinitive reports whether err settles the discovery request, so that
+// retrying it cannot change the outcome. Only the caller decides what a
+// definitive failure means: NotFound is absence, Forbidden is a
misconfiguration.
+func isDefinitive(err error) bool {
+ switch {
+ // Discovery of a group/version that is not installed answers 404.
+ case apierrors.IsNotFound(err):
+ return true
+ // RBAC forbids discovery: waiting will not change that, but it says
nothing
+ // about whether the resource is served.
+ case apierrors.IsForbidden(err):
+ return true
+ // Connection refused, timeouts, EOF, 5xx, throttling: the API server
may
+ // answer differently once it is reachable again, so no conclusion yet.
+ default:
+ return false
+ }
}
func FormatGVK(obj client.Object) string {
diff --git a/pkg/utils/k8s_test.go b/pkg/utils/k8s_test.go
new file mode 100644
index 00000000..73ce3579
--- /dev/null
+++ b/pkg/utils/k8s_test.go
@@ -0,0 +1,324 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you 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 utils
+
+import (
+ "context"
+ "errors"
+ "net"
+ "net/http"
+ "net/http/httptest"
+ "testing"
+ "time"
+
+ "github.com/go-logr/logr"
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+ apierrors "k8s.io/apimachinery/pkg/api/errors"
+ "k8s.io/apimachinery/pkg/runtime/schema"
+ "k8s.io/apimachinery/pkg/util/wait"
+ "k8s.io/client-go/discovery"
+ "k8s.io/client-go/rest"
+)
+
+// fastBackoff keeps the tests quick while preserving the shape of the
+// production backoffs (a growing interval over a fixed number of steps).
+func fastBackoff(steps int) wait.Backoff {
+ return wait.Backoff{Duration: time.Millisecond, Factor: backoffFactor,
Jitter: backoffJitter, Steps: steps}
+}
+
+// declaredBudget is how long a backoff sleeps in total before jitter, i.e. the
+// sum of the intervals between its attempts.
+func declaredBudget(backoff wait.Backoff) time.Duration {
+ var total, interval time.Duration = 0, backoff.Duration
+ for i := 0; i < backoff.Steps-1; i++ {
+ total += interval
+ interval = time.Duration(float64(interval) * backoff.Factor)
+ }
+ return total
+}
+
+// countingBackoff reports how many times a backoff actually runs its
condition.
+func countAttempts(backoff wait.Backoff) int {
+ var attempts int
+ _ = wait.ExponentialBackoff(backoff, func() (bool, error) {
+ attempts++
+ return false, nil
+ })
+ return attempts
+}
+
+// TestBackoffsUseAllTheirSteps guards against a subtle wait.Backoff behavior:
+// setting Cap zeroes the remaining steps as soon as the capped interval is
+// reached, silently cutting the retries short. Both backoffs must run every
+// step they declare.
+func TestBackoffsUseAllTheirSteps(t *testing.T) {
+ for _, tt := range []struct {
+ name string
+ backoff wait.Backoff
+ steps int
+ budget time.Duration // as documented next to the constants
+ }{
+ {"discovery", discoveryBackoff(), discoveryRetrySteps, 7500 *
time.Millisecond},
+ {"apiServerWait", apiServerWaitBackoff(), apiServerWaitSteps,
31 * time.Second},
+ } {
+ t.Run(tt.name, func(t *testing.T) {
+ assert.Zero(t, tt.backoff.Cap, "Cap truncates the
retries, see wait.delay()")
+ assert.Equal(t, tt.budget, declaredBudget(tt.backoff),
"the documented budget must match the backoff")
+
+ b := tt.backoff
+ b.Duration = time.Microsecond // keep the test fast,
preserve the shape
+ assert.Equal(t, tt.steps, countAttempts(b))
+ })
+ }
+}
+
+func TestIsDefinitive(t *testing.T) {
+ gr := schema.GroupResource{Group: "gateway.networking.k8s.io",
Resource: "tcproutes"}
+ tests := []struct {
+ name string
+ err error
+ definitive bool
+ }{
+ // Definitive: the group/version is genuinely not served.
+ {"not found", apierrors.NewNotFound(gr, "tcproutes"), true},
+ // Definitive: RBAC will not self-heal by waiting.
+ {"forbidden", apierrors.NewForbidden(gr, "tcproutes",
errors.New("nope")), true},
+ // Indeterminate: the API server is temporarily unreachable or
overloaded.
+ {"connection refused", &net.OpError{Op: "dial", Err:
errors.New("connection refused")}, false},
+ {"service unavailable",
apierrors.NewServiceUnavailable("apiserver is starting up"), false},
+ {"server timeout", apierrors.NewServerTimeout(gr, "get", 1),
false},
+ {"too many requests", apierrors.NewTooManyRequestsError("slow
down"), false},
+ {"internal error",
apierrors.NewInternalError(errors.New("boom")), false},
+ {"generic transport error", errors.New("EOF"), false},
+ }
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ assert.Equal(t, tt.definitive, isDefinitive(tt.err))
+ })
+ }
+}
+
+// TestRetryUntilDefinitive_ReportsIndeterminateFailure is the core regression
+// test for #2734: a discovery failure that gives no definitive answer must
+// surface as an error so the caller aborts startup, instead of collapsing into
+// "resource absent" and permanently skipping the field index and its
controller.
+func TestRetryUntilDefinitive_ReportsIndeterminateFailure(t *testing.T) {
+ var calls int
+ probe := func() error {
+ calls++
+ return apierrors.NewServiceUnavailable("apiserver not ready")
+ }
+
+ err := retryUntilDefinitive(fastBackoff(3), logr.Discard(), probe)
+ require.Error(t, err, "an unreachable API server must not be reported
as a definitive answer")
+ assert.False(t, isDefinitive(err))
+ assert.Equal(t, 3, calls, "the whole retry budget must be used")
+}
+
+func TestRetryUntilDefinitive_TransientThenSuccess(t *testing.T) {
+ var calls int
+ probe := func() error {
+ calls++
+ if calls < 3 {
+ return apierrors.NewServiceUnavailable("apiserver not
ready")
+ }
+ return nil
+ }
+
+ require.NoError(t, retryUntilDefinitive(fastBackoff(5), logr.Discard(),
probe))
+ assert.Equal(t, 3, calls, "an indeterminate failure must be retried,
not given up on")
+}
+
+// TestRetryUntilDefinitive_DefinitiveAbsent ensures a genuinely absent CRD is
+// answered immediately without retrying, preserving optional-CRD support.
+func TestRetryUntilDefinitive_DefinitiveAbsent(t *testing.T) {
+ gr := schema.GroupResource{Group: "gateway.networking.k8s.io",
Resource: "tcproutes"}
+ var calls int
+ probe := func() error {
+ calls++
+ return apierrors.NewNotFound(gr, "tcproutes")
+ }
+
+ err := retryUntilDefinitive(fastBackoff(5), logr.Discard(), probe)
+ assert.True(t, isDefinitive(err))
+ assert.Equal(t, 1, calls, "a definitive NotFound must not be retried")
+}
+
+// TestRetryUntilDefinitive_NeverProbed guards the caller's assumption that a
+// nil error means "the probe succeeded": a backoff with no steps left never
runs
+// the condition, and reporting nil there would be read as "resource present".
+func TestRetryUntilDefinitive_NeverProbed(t *testing.T) {
+ var calls int
+ probe := func() error {
+ calls++
+ return nil
+ }
+
+ err := retryUntilDefinitive(wait.Backoff{Steps: 0}, logr.Discard(),
probe)
+ require.Error(t, err)
+ assert.Zero(t, calls)
+ assert.False(t, isDefinitive(err), "an unattempted probe must not look
like a definitive answer either")
+}
+
+func TestWaitForAPIServer_RetriesUntilReachable(t *testing.T) {
+ var calls int
+ probe := func() error {
+ calls++
+ if calls < 3 {
+ return &net.OpError{Op: "dial", Err:
errors.New("connection refused")}
+ }
+ return nil
+ }
+
+ require.NoError(t, waitForAPIServer(context.Background(), probe,
fastBackoff(6), logr.Discard()))
+ assert.Equal(t, 3, calls)
+}
+
+// TestWaitForAPIServer_AuthRejectionIsReachable covers a hardened cluster that
+// does not grant the controller access to the probed endpoint: the API server
+// answered, so startup must proceed instead of stalling until the budget runs
+// out and then failing on a perfectly healthy cluster.
+func TestWaitForAPIServer_AuthRejectionIsReachable(t *testing.T) {
+ for name, probeErr := range map[string]error{
+ "unauthorized": apierrors.NewUnauthorized("no token"),
+ "forbidden": apierrors.NewForbidden(schema.GroupResource{},
"", errors.New("nope")),
+ } {
+ t.Run(name, func(t *testing.T) {
+ var calls int
+ probe := func() error {
+ calls++
+ return probeErr
+ }
+ require.NoError(t,
waitForAPIServer(context.Background(), probe, fastBackoff(5), logr.Discard()))
+ assert.Equal(t, 1, calls, "an answer, even a rejection,
must not be retried")
+ })
+ }
+}
+
+func TestWaitForAPIServer_GivesUpWithTheLastError(t *testing.T) {
+ probe := func() error { return
apierrors.NewServiceUnavailable("apiserver not ready") }
+
+ err := waitForAPIServer(context.Background(), probe, fastBackoff(3),
logr.Discard())
+ require.Error(t, err)
+ assert.True(t, apierrors.IsServiceUnavailable(errors.Unwrap(err)), "the
last probe error must be reported, got %v", err)
+}
+
+// TestWaitForAPIServer_HonorsContext ensures the wait is interrupted
mid-flight,
+// so a shutdown signal is not ignored while the API server is unreachable.
+func TestWaitForAPIServer_HonorsContext(t *testing.T) {
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
+
+ var calls int
+ probe := func() error {
+ if calls++; calls == 2 {
+ cancel()
+ }
+ return &net.OpError{Op: "dial", Err: errors.New("connection
refused")}
+ }
+
+ // A budget far larger than the number of attempts the test expects:
the wait
+ // must end because ctx was cancelled, not because it ran out of steps.
+ err := waitForAPIServer(ctx, probe, fastBackoff(100), logr.Discard())
+ require.ErrorIs(t, err, context.Canceled)
+ assert.Equal(t, 2, calls, "the probe must stop being called once ctx is
cancelled")
+}
+
+// newDiscoveryClient serves discovery responses from handler.
+func newDiscoveryClient(t *testing.T, handler http.HandlerFunc)
discovery.DiscoveryInterface {
+ t.Helper()
+ srv := httptest.NewServer(handler)
+ t.Cleanup(srv.Close)
+ c, err := discovery.NewDiscoveryClientForConfig(&rest.Config{Host:
srv.URL})
+ require.NoError(t, err)
+ return c
+}
+
+// TestHasAPIResource covers the contract every caller depends on, against a
real
+// discovery client rather than a fake: only a definitive answer may resolve to
+// "absent", everything else must be an error.
+func TestHasAPIResource(t *testing.T) {
+ gvk := schema.GroupVersionKind{Group: "gateway.networking.k8s.io",
Version: "v1alpha2", Kind: "TCPRoute"}
+
+ tests := []struct {
+ name string
+ handler http.HandlerFunc
+ wantFound bool
+ wantErr bool
+ }{
+ {
+ name: "kind served",
+ handler: func(w http.ResponseWriter, _ *http.Request) {
+ w.Header().Set("Content-Type",
"application/json")
+ _, _ =
w.Write([]byte(`{"kind":"APIResourceList","groupVersion":"gateway.networking.k8s.io/v1alpha2",`
+
+
`"resources":[{"name":"tcproutes","kind":"TCPRoute"}]}`))
+ },
+ wantFound: true,
+ },
+ {
+ name: "group/version served but kind missing",
+ handler: func(w http.ResponseWriter, _ *http.Request) {
+ w.Header().Set("Content-Type",
"application/json")
+ _, _ =
w.Write([]byte(`{"kind":"APIResourceList","groupVersion":"gateway.networking.k8s.io/v1alpha2",`
+
+
`"resources":[{"name":"udproutes","kind":"UDPRoute"}]}`))
+ },
+ },
+ {
+ // What a real API server returns for a group/version
that is not
+ // installed: a plain-text 404, not a Status object.
+ name: "group/version not installed",
+ handler: func(w http.ResponseWriter, _ *http.Request) {
+ http.NotFound(w, nil)
+ },
+ },
+ {
+ // The #2734 scenario: the API server cannot answer.
This must not be
+ // reported as "the CRD is not installed".
+ name: "api server unavailable",
+ handler: func(w http.ResponseWriter, _ *http.Request) {
+ w.WriteHeader(http.StatusServiceUnavailable)
+ },
+ wantErr: true,
+ },
+ {
+ // Forbidden is definitive but proves nothing about the
resource: the
+ // CRD may well be served and only discovery denied.
Resolving it to
+ // "absent" would permanently skip the controller, so
it must abort
+ // startup and point at the RBAC instead.
+ name: "discovery forbidden",
+ handler: func(w http.ResponseWriter, _ *http.Request) {
+ w.WriteHeader(http.StatusForbidden)
+ },
+ wantErr: true,
+ },
+ }
+
+ for _, tt := range tests {
+ t.Run(tt.name, func(t *testing.T) {
+ found, err := hasAPIResource(newDiscoveryClient(t,
tt.handler), gvk, fastBackoff(2), logr.Discard())
+ if tt.wantErr {
+ require.Error(t, err)
+ assert.False(t, found)
+ return
+ }
+ require.NoError(t, err)
+ assert.Equal(t, tt.wantFound, found)
+ })
+ }
+}