Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 6 additions & 4 deletions cmd/controller/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,11 @@ func main() {
klog.Fatal("Config not found")
}

// Configure client side rate limiting settings
// @TODO: make this configurable
config.QPS = 20
config.Burst = 30

leaseLockNamespace := util.GetNamespace()
leaseLockId := uuid.New().String()

Expand Down Expand Up @@ -174,10 +179,7 @@ func getDefaultConcurrencyConfig() map[int]int {
if !ok {
defaultReconcileForResource = controller.DefaultReconcile
}
concurrencyValue := getConcurrencyConfigForResource(resourceEnvSuffix, defaultReconcileForResource)
if concurrencyValue < 1 {
concurrencyValue = 1
}
concurrencyValue := max(getConcurrencyConfigForResource(resourceEnvSuffix, defaultReconcileForResource), 1)
concurrencyConfig[resourceKey] = concurrencyValue
}
return concurrencyConfig
Expand Down
49 changes: 48 additions & 1 deletion internal/controller/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,13 @@ SPDX-License-Identifier: Apache-2.0
package controller

import (
"context"
"net/url"
"os"
"time"

"github.com/prometheus/client_golang/prometheus"
k8sclientmetrics "k8s.io/client-go/tools/metrics"
"k8s.io/client-go/util/workqueue"
)

Expand Down Expand Up @@ -134,10 +138,47 @@ var (
Name: Retries,
Help: "Retries in workqueue",
}, []string{"name"})

// K8s client-go metrics aren't exposed. This is a copy of: https://github.com/kubernetes/kubernetes/blob/master/staging/src/k8s.io/component-base/metrics/prometheus/restclient/metrics.go#L78
rateLimiterLatency = prometheus.NewHistogramVec(prometheus.HistogramOpts{
Namespace: CAPOp,
Name: "rest_client_rate_limiter_duration_seconds",
Help: "Client side rate limiter latency in seconds. Broken down by verb, and host.",
Buckets: []float64{0.005, 0.025, 0.1, 0.25, 0.5, 1.0, 2.0, 4.0, 8.0, 15.0, 30.0, 60.0},
},
[]string{"verb", "host"},
)

// K8s client-go metrics aren't exposed. This is a copy of: https://github.com/kubernetes/kubernetes/blob/master/staging/src/k8s.io/component-base/metrics/prometheus/restclient/metrics.go#L88
requestResult = prometheus.NewCounterVec(prometheus.CounterOpts{
Namespace: CAPOp,
Name: "rest_client_requests_total",
Help: "Number of HTTP requests, partitioned by status code, method, and host.",
}, []string{"code", "method", "host"})
)

// #region k8sRequestResultProvider
// This isn't exposed by K8s, so we made a copy of: https://github.com/kubernetes/kubernetes/blob/master/staging/src/k8s.io/component-base/metrics/prometheus/restclient/metrics.go#L252
type k8sRequestResultProvider struct {
m *prometheus.CounterVec
}

func (p *k8sRequestResultProvider) Increment(_ context.Context, code string, method string, host string) {
p.m.WithLabelValues(code, method, host).Inc()
}

type k8sRequestlatencyAdapter struct {
m *prometheus.HistogramVec
}

func (l *k8sRequestlatencyAdapter) Observe(_ context.Context, verb string, u url.URL, latency time.Duration) {
l.m.WithLabelValues(verb, u.Host).Observe(latency.Seconds())
}

// #endregion

// Create a variable to hold all the collectors
var collectors = []prometheus.Collector{ReconcileErrors, Panics, TenantOperations, ServiceOperations, depth, adds, latency, workDuration, unfinished, longestRunningProcessor, retries}
var collectors = []prometheus.Collector{ReconcileErrors, Panics, TenantOperations, ServiceOperations, depth, adds, latency, workDuration, unfinished, longestRunningProcessor, retries, requestResult}

// #region capOperatorMetricsProvider
// capOperatorMetricsProvider implements workqueue.MetricsProvider
Expand Down Expand Up @@ -184,6 +225,12 @@ func initializeMetrics() {
// Register CAP Operator metrics
prometheus.MustRegister(collectors...)

// Register Kubernetes client-go REST API metrics with the custom k8sRequestResultProvider
k8sclientmetrics.Register(k8sclientmetrics.RegisterOpts{
RequestResult: &k8sRequestResultProvider{requestResult},
RateLimiterLatency: &k8sRequestlatencyAdapter{rateLimiterLatency},
})

// Register CAP Operator metrics provider as the workqueue metrics provider (needed for the workqueue metrics, to be done just once)
workqueue.SetProvider(capOperatorMetricsProvider{})
}
Expand Down
115 changes: 48 additions & 67 deletions internal/controller/reconcile-capapplication.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ import (

"github.com/sap/cap-operator/internal/util"
"github.com/sap/cap-operator/pkg/apis/sme.sap.com/v1alpha1"
"golang.org/x/sync/errgroup"
corev1 "k8s.io/api/core/v1"
k8sErrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
Expand All @@ -39,8 +38,7 @@ const (
)

func (c *Controller) reconcileCAPApplication(ctx context.Context, item QueueItem, _ int) (result *ReconcileResult, err error) {
lister := c.crdInformerFactory.Sme().V1alpha1().CAPApplications().Lister()
cached, err := lister.CAPApplications(item.ResourceKey.Namespace).Get(item.ResourceKey.Name)
cached, err := c.crdInformerFactory.Sme().V1alpha1().CAPApplications().Lister().CAPApplications(item.ResourceKey.Namespace).Get(item.ResourceKey.Name)
if err != nil {
return nil, handleOperatorResourceErrors(err)
}
Expand Down Expand Up @@ -78,7 +76,8 @@ func (c *Controller) reconcileCAPApplication(ctx context.Context, item QueueItem
result, err = c.handleCAPApplicationDependentResources(ctx, ca)
}

return c.checkAdditionalConditions(ca, result, err)
result, err = c.checkAdditionalConditions(ctx, ca, result, err)
return
}

func (c *Controller) handleCAPApplicationDependentResources(ctx context.Context, ca *v1alpha1.CAPApplication) (requeue *ReconcileResult, err error) {
Expand Down Expand Up @@ -131,15 +130,10 @@ func (c *Controller) handleCAPApplicationDependentResources(ctx context.Context,
}

// step 5 - reconcile service exposure, create/update services based on the latest CAV
if err = c.reconcileServiceNetworking(ctx, ca, cav); err != nil {
return
}

// step 6 - check and set consistent status; check for newer versions and trigger tenant networking updates
return c.verifyApplicationConsistent(ctx, ca)
return nil, c.reconcileServiceNetworking(ctx, ca, cav)
}

func (c *Controller) verifyApplicationConsistent(ctx context.Context, ca *v1alpha1.CAPApplication) (requeue *ReconcileResult, err error) {
func (c *Controller) verifyApplicationConsistent(ca *v1alpha1.CAPApplication) {
if ca.Status.State != v1alpha1.CAPApplicationStateConsistent {
ca.SetStatusWithReadyCondition(v1alpha1.CAPApplicationStateConsistent, metav1.ConditionTrue, "VersionExists", "")
// Update additional condition `LatestVersionReady` to True
Expand All @@ -150,56 +144,59 @@ func (c *Controller) verifyApplicationConsistent(ctx context.Context, ca *v1alph
ca.SetStatusCondition(string(v1alpha1.ConditionTypeAllTenantsReady), metav1.ConditionTrue, "ProviderTenantReady", "")
}
}

// Check for newer CAPApplicationVersion and trigger tenant networking updates
return nil, c.checkNewCavAndTenantNetworking(ctx, ca)
}

func (c *Controller) checkNewCavAndTenantNetworking(ctx context.Context, ca *v1alpha1.CAPApplication) error {
// Get the latest CAV for the tenant
cav, err := c.getLatestReadyCAPApplicationVersion(ca, false)
if err != nil {
return err
func (c *Controller) checkNewCavAndTenantReconcile(ctx context.Context, ca *v1alpha1.CAPApplication, readyCav *v1alpha1.CAPApplicationVersion, latestCav *v1alpha1.CAPApplicationVersion) (*ReconcileResult, error) {
// No tenants for services only scenario
if ca.IsServicesOnly() {
return nil, nil
}

// Reset ready Condition and Reason for Tenant check AllTenantsReady --> True
readyCondition := metav1.ConditionTrue
readyReason := string(v1alpha1.ConditionTypeAllTenantsReady)

// Get all relevant tenants
tenants, err := c.getRelevantTenantsForCA(ca)
if err != nil || len(tenants) == 0 {
return err
return nil, err
}

netUpdGrp := errgroup.Group{}
checkDone := false
updated := false
var result *ReconcileResult
if ca.Annotations[AnnotationEnableVersionAffinity] == "true" {
result = &ReconcileResult{}
}
for _, tenant := range tenants {
if tenant.Status.CurrentCAPApplicationVersionInstance != "" {
t := tenant
netUpdGrp.Go(func() error {
return c.reconcileTenantNetworking(ctx, t, t.Status.CurrentCAPApplicationVersionInstance, ca)
})
}

if upd, err := c.checkForTenantVersionUpgrade(ctx, ca, cav, tenant); err != nil {
return err
if upd, err := c.checkForTenantVersionUpgradeAndReconcile(ctx, ca, readyCav, tenant, result); err != nil {
return result, err
} else if upd {
updated = true
}
// When a Tenant state is not Ready -or- when version of tenant (with VersionUpgradeStrategy = always) does not match the latest CAV version --> AllTenantsReady = False
if !checkDone && (updated || (tenant.Status.State != v1alpha1.CAPTenantStateReady || (tenant.Spec.VersionUpgradeStrategy == v1alpha1.VersionUpgradeStrategyTypeAlways && latestCav.Spec.Version != tenant.Spec.Version))) {
readyCondition = metav1.ConditionFalse
readyReason = "NotAllTenantsReady"
checkDone = true
}
}

if err = netUpdGrp.Wait(); err != nil {
return fmt.Errorf("failed to reconcile tenant networking: %w", err)
}
// Update `AllTenantsReady` status condition
ca.SetStatusCondition(string(v1alpha1.ConditionTypeAllTenantsReady), readyCondition, readyReason, "")

if updated {
msg := fmt.Sprintf("new version %s.%s was used to trigger tenant upgrades", cav.Namespace, cav.Name)
msg := fmt.Sprintf("new version %s.%s was used to trigger tenant upgrades", readyCav.Namespace, readyCav.Name)
ca.SetStatusWithReadyCondition(v1alpha1.CAPApplicationStateProcessing, metav1.ConditionFalse, CAPApplicationEventNewCAVTriggeredTenantUpgrade, msg)
ca.SetStatusCondition(string(v1alpha1.ConditionTypeLatestVersionReady), metav1.ConditionTrue, string(v1alpha1.ConditionTypeLatestVersionReady), "")
ca.SetStatusCondition(string(v1alpha1.ConditionTypeAllTenantsReady), metav1.ConditionFalse, "UpgradingTenants", "")
c.Event(ca, nil, corev1.EventTypeNormal, CAPApplicationEventNewCAVTriggeredTenantUpgrade, EventActionCheckForVersion, msg)
}
return nil
return result, nil
}

func (c *Controller) checkForTenantVersionUpgrade(ctx context.Context, ca *v1alpha1.CAPApplication, cav *v1alpha1.CAPApplicationVersion, tenant *v1alpha1.CAPTenant) (bool, error) {
func (c *Controller) checkForTenantVersionUpgradeAndReconcile(ctx context.Context, ca *v1alpha1.CAPApplication, cav *v1alpha1.CAPApplicationVersion, tenant *v1alpha1.CAPTenant, result *ReconcileResult) (bool, error) {
// This is done to reconcile tenant networking which may be needed for session affinity
if result != nil && tenant.Status.CurrentCAPApplicationVersionInstance != "" {
result.AddResource(ResourceCAPTenant, tenant.Name, tenant.Namespace, 1)
}

if tenant.Spec.VersionUpgradeStrategy == v1alpha1.VersionUpgradeStrategyTypeNever {
// Skip non relevant tenants
return false, nil
Expand All @@ -226,57 +223,41 @@ func (c *Controller) checkForTenantVersionUpgrade(ctx context.Context, ca *v1alp
return false, nil
}

func (c *Controller) checkAdditionalConditions(ca *v1alpha1.CAPApplication, result *ReconcileResult, err error) (*ReconcileResult, error) {
func (c *Controller) checkAdditionalConditions(ctx context.Context, ca *v1alpha1.CAPApplication, result *ReconcileResult, err error) (*ReconcileResult, error) {
// In case of explicit Reconcile or errors return back with the original result
if result != nil || err != nil {
return result, err
}

// check and set consistent status;
c.verifyApplicationConsistent(ca)

// Check and update additional status conditions
// Set ready Condition and Reason for Version check LatestVersionNotReady = True
readyCondition := metav1.ConditionTrue
readyReason := string(v1alpha1.ConditionTypeLatestVersionReady)

// Get latest CAV (incl. ones that may not be ready)
cav, err := c.getLatestCAPApplicationVersion(ca)
readyCav := cav
if err != nil {
return nil, err
}
// When the latest CAV is not Ready --> LatestVersionNotReady = False
if cav.Status.State != v1alpha1.CAPApplicationVersionStateReady {
readyCondition = metav1.ConditionFalse
readyReason = "LatestVersionNotReady"
// Get the latest Ready CAV for the tenant
readyCav, err = c.getLatestReadyCAPApplicationVersion(ca, false)
if err != nil {
return nil, err
}
}

// Update `LatestVersionReady` status condition
ca.SetStatusCondition(string(v1alpha1.ConditionTypeLatestVersionReady), readyCondition, readyReason, "")

// No tenants for services only scenario
if ca.IsServicesOnly() {
return nil, nil
}

// Reset ready Condition and Reason for Tenant check AllTenantsReady --> True
readyCondition = metav1.ConditionTrue
readyReason = string(v1alpha1.ConditionTypeAllTenantsReady)

// Get all relevant tenants
tenants, err := c.getRelevantTenantsForCA(ca)
if err != nil {
return nil, err
}
for _, tenant := range tenants {
// When a Tenant state is not Ready -or- when version of tenant (with VersionUpgradeStrategy = always) does not match the latest CAV version --> AllTenantsReady = False
if tenant.Status.State != v1alpha1.CAPTenantStateReady || (tenant.Spec.VersionUpgradeStrategy == v1alpha1.VersionUpgradeStrategyTypeAlways && cav.Spec.Version != tenant.Spec.Version) {
readyCondition = metav1.ConditionFalse
readyReason = "NotAllTenantsReady"
break
}
}
// Update `AllTenantsReady` status condition
ca.SetStatusCondition(string(v1alpha1.ConditionTypeAllTenantsReady), readyCondition, readyReason, "")

return nil, nil
return c.checkNewCavAndTenantReconcile(ctx, ca, readyCav, cav)
}

func (c *Controller) updateCAPApplication(ctx context.Context, ca *v1alpha1.CAPApplication) error {
Expand Down
Loading
Loading