Skip to content
Open
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
33 changes: 30 additions & 3 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import (
// to ensure that exec-entrypoint and run can make use of them.
"github.com/KimMachineGun/automemlimit/memlimit"
"golang.org/x/sync/errgroup"
autoscalingv2 "k8s.io/api/autoscaling/v2"
_ "k8s.io/client-go/plugin/pkg/client/auth"

apimeta "k8s.io/apimachinery/pkg/api/meta"
Expand All @@ -31,6 +32,7 @@ import (
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/cache"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/client/apiutil"
"sigs.k8s.io/controller-runtime/pkg/cluster"
"sigs.k8s.io/controller-runtime/pkg/healthz"
"sigs.k8s.io/controller-runtime/pkg/log/zap"
Expand Down Expand Up @@ -97,6 +99,19 @@ func init() {
// +kubebuilder:scaffold:scheme
}

func managedResourceGVKs(s *runtime.Scheme, objs ...client.Object) ([]schema.GroupVersionKind, error) {
gvks := make([]schema.GroupVersionKind, 0, len(objs))
for _, obj := range objs {
gvk, err := apiutil.GVKForObject(obj, s)
if err != nil {
return nil, err
}
gvks = append(gvks, gvk)
}

return gvks, nil
}

//nolint:gocyclo // main wires all controller paths; complexity is inherent to startup sequencing
func main() {

Expand Down Expand Up @@ -361,6 +376,11 @@ func main() {
setupLog.Error(err, "unable to create controller", "controller", "WorkloadDeployment")
os.Exit(1)
}

if err = (&controller.WorkloadDeploymentHPAReconciler{}).SetupWithManager(mgr); err != nil {
setupLog.Error(err, "unable to create controller", "controller", "WorkloadDeploymentHPA")
os.Exit(1)
}
}

if enableCellControllers {
Expand Down Expand Up @@ -560,6 +580,15 @@ func initializeClusterDiscovery(
return nil, nil, "", nil, fmt.Errorf("unable to create root client for service-catalog: %w", err)
}

managedResources, err := managedResourceGVKs(
scheme,
&computev1alpha.Instance{},
&autoscalingv2.HorizontalPodAutoscaler{},
)
if err != nil {
return nil, nil, "", nil, fmt.Errorf("unable to resolve managed resource GVKs: %w", err)
}

provider, err = consumerprovider.New(providerMgr, consumerprovider.Options{
RootClient: rootClient,
Scheme: scheme,
Expand All @@ -570,9 +599,7 @@ func initializeClusterDiscovery(
o.Cache.DefaultTransform = cache.TransformStripManagedFields()
},
},
ManagedResources: []schema.GroupVersionKind{
computev1alpha.GroupVersion.WithKind("Instance"),
},
ManagedResources: managedResources,
Teardowns: []consumerprovider.Teardown{
controller.NewComputeTeardown(quotaClientManager, federationClient, scheme),
},
Expand Down
12 changes: 12 additions & 0 deletions config/components/controller_rbac/role.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,18 @@ rules:
verbs:
- get
- list
- apiGroups:
- autoscaling
resources:
- horizontalpodautoscalers
verbs:
- create
- delete
- get
- list
- patch
- update
- watch
- apiGroups:
- compute.datumapis.com
resources:
Expand Down
12 changes: 8 additions & 4 deletions internal/controller/teardown.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,14 @@ import (
quotametrics "go.datum.net/compute/internal/quota"
)

// labelServiceName is the label key the consumer provider uses to scope
// deactivation cleanup. Every Instance the compute operator creates in a
// consumer project carries this label so disengage can target them by service.
const labelServiceName = "services.miloapis.com/service-name"
// labelServiceName and labelServiceValue are the label key and value the
// consumer provider uses to scope deactivation cleanup. Every resource the
// compute operator creates in a consumer project carries this label so disengage
// can target them by service.
const (
labelServiceName = "services.miloapis.com/service-name"
labelServiceValue = "compute.datumapis.com"
)

// ComputeTeardown implements consumer.Teardown. It is invoked after the
// consumer provider has cancelled the per-cluster context and marked labeled
Expand Down
2 changes: 2 additions & 0 deletions internal/controller/testing_helpers_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"fmt"
"sync"

autoscalingv2 "k8s.io/api/autoscaling/v2"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/client-go/tools/events"
Expand All @@ -26,6 +27,7 @@ import (
// cluster (corev1 + compute).
func newProjectScheme() *runtime.Scheme {
s := runtime.NewScheme()
_ = autoscalingv2.AddToScheme(s)
_ = corev1.AddToScheme(s)
_ = computev1alpha.AddToScheme(s)
return s
Expand Down
207 changes: 207 additions & 0 deletions internal/controller/workloaddeployment_hpa_controller.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,207 @@
// SPDX-License-Identifier: AGPL-3.0-only

package controller

import (
"context"
"fmt"

autoscalingv2 "k8s.io/api/autoscaling/v2"
corev1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/controller-runtime/pkg/log"
"sigs.k8s.io/controller-runtime/pkg/predicate"
mcbuilder "sigs.k8s.io/multicluster-runtime/pkg/builder"
mccontext "sigs.k8s.io/multicluster-runtime/pkg/context"
mcmanager "sigs.k8s.io/multicluster-runtime/pkg/manager"
mcreconcile "sigs.k8s.io/multicluster-runtime/pkg/reconcile"

computev1alpha "go.datum.net/compute/api/v1alpha"
)

// WorkloadDeploymentHPAReconciler manages cell-local HorizontalPodAutoscalers
// for WorkloadDeployments that opt into load-driven autoscaling.
type WorkloadDeploymentHPAReconciler struct {
mgr mcmanager.Manager
}

// +kubebuilder:rbac:groups=compute.datumapis.com,resources=workloaddeployments,verbs=get;list;watch
// +kubebuilder:rbac:groups=autoscaling,resources=horizontalpodautoscalers,verbs=get;list;watch;create;update;patch;delete

func (r *WorkloadDeploymentHPAReconciler) Reconcile(ctx context.Context, req mcreconcile.Request) (ctrl.Result, error) {
logger := log.FromContext(ctx)

cl, err := r.mgr.GetCluster(ctx, req.ClusterName)
if err != nil {
return ctrl.Result{}, err
}

ctx = mccontext.WithCluster(ctx, req.ClusterName)

var deployment computev1alpha.WorkloadDeployment
if err := cl.GetClient().Get(ctx, req.NamespacedName, &deployment); err != nil {
return ctrl.Result{}, client.IgnoreNotFound(err)
}

if !deployment.DeletionTimestamp.IsZero() {
return ctrl.Result{}, nil
}

if !workloadDeploymentAutoscalingEnabled(&deployment) {
return ctrl.Result{}, deleteWorkloadDeploymentHPA(ctx, cl.GetClient(), &deployment)
}

logger.Info("reconciling deployment HPA")
defer logger.Info("deployment HPA reconcile complete")

metrics, err := workloadDeploymentHPAMetrics(&deployment)
if err != nil {
return ctrl.Result{}, err
}

hpa := autoscalingv2.HorizontalPodAutoscaler{
ObjectMeta: metav1.ObjectMeta{
Name: deployment.Name,
Namespace: deployment.Namespace,
},
}

_, err = controllerutil.CreateOrPatch(ctx, cl.GetClient(), &hpa, func() error {
if hpa.UID != "" && !metav1.IsControlledBy(&hpa, &deployment) {
return fmt.Errorf("HPA %s/%s already exists and is not controlled by WorkloadDeployment %s/%s",
hpa.Namespace, hpa.Name, deployment.Namespace, deployment.Name)
}

if err := controllerutil.SetControllerReference(&deployment, &hpa, cl.GetScheme()); err != nil {
return err
}

hpa.Labels = workloadDeploymentHPALabels(&deployment)

hpa.Spec = autoscalingv2.HorizontalPodAutoscalerSpec{
ScaleTargetRef: autoscalingv2.CrossVersionObjectReference{
APIVersion: computev1alpha.GroupVersion.String(),
Kind: "WorkloadDeployment",
Name: deployment.Name,
},
MinReplicas: new(deployment.Spec.ScaleSettings.MinReplicas),
MaxReplicas: *deployment.Spec.ScaleSettings.MaxReplicas,
Metrics: metrics,
}

return nil
})
if err != nil {
return ctrl.Result{}, fmt.Errorf("failed reconciling deployment HPA: %w", err)
}

return ctrl.Result{}, nil
}

func workloadDeploymentAutoscalingEnabled(deployment *computev1alpha.WorkloadDeployment) bool {
return deployment.Spec.ScaleSettings.MaxReplicas != nil && len(deployment.Spec.ScaleSettings.Metrics) > 0
}

func deleteWorkloadDeploymentHPA(ctx context.Context, c client.Client, deployment *computev1alpha.WorkloadDeployment) error {
var hpa autoscalingv2.HorizontalPodAutoscaler
if err := c.Get(ctx, client.ObjectKey{Namespace: deployment.Namespace, Name: deployment.Name}, &hpa); err != nil {
if apierrors.IsNotFound(err) {
return nil
}
return fmt.Errorf("failed fetching deployment HPA: %w", err)
}
if !metav1.IsControlledBy(&hpa, deployment) {
return nil
}

if err := c.Delete(ctx, &hpa); client.IgnoreNotFound(err) != nil {
return fmt.Errorf("failed deleting deployment HPA: %w", err)
}

return nil
}

func workloadDeploymentHPALabels(deployment *computev1alpha.WorkloadDeployment) map[string]string {
return map[string]string{
labelServiceName: labelServiceValue,
computev1alpha.WorkloadDeploymentUIDLabel: string(deployment.UID),
computev1alpha.WorkloadDeploymentNameLabel: deployment.Name,
computev1alpha.WorkloadNameLabel: deployment.Spec.WorkloadRef.Name,
computev1alpha.PlacementNameLabel: deployment.Spec.PlacementName,
computev1alpha.CityCodeLabel: deployment.Spec.CityCode,
}
}

func workloadDeploymentHPAMetrics(deployment *computev1alpha.WorkloadDeployment) ([]autoscalingv2.MetricSpec, error) {
metrics := make([]autoscalingv2.MetricSpec, 0, len(deployment.Spec.ScaleSettings.Metrics))
for i, metric := range deployment.Spec.ScaleSettings.Metrics {
if metric.Resource == nil {
return nil, fmt.Errorf("metric %d has no resource source", i)
}

if metric.Resource.Name != corev1.ResourceCPU && metric.Resource.Name != corev1.ResourceMemory {
return nil, fmt.Errorf("metric %d uses unsupported resource %q", i, metric.Resource.Name)
}

target, err := workloadDeploymentHPAMetricTarget(metric.Resource.Target)
if err != nil {
return nil, fmt.Errorf("metric %d has invalid target: %w", i, err)
}

metrics = append(metrics, autoscalingv2.MetricSpec{
Type: autoscalingv2.ResourceMetricSourceType,
Resource: &autoscalingv2.ResourceMetricSource{
Name: metric.Resource.Name,
Target: target,
},
})
}

return metrics, nil
}

func workloadDeploymentHPAMetricTarget(target computev1alpha.MetricTarget) (autoscalingv2.MetricTarget, error) {
setTargets := 0
if target.Value != nil {
setTargets++
}
if target.AverageValue != nil {
setTargets++
}
if target.AverageUtilization != nil {
setTargets++
}
if setTargets != 1 {
return autoscalingv2.MetricTarget{}, fmt.Errorf("exactly one target value must be set")
}

if target.Value != nil {
value := target.Value.DeepCopy()
return autoscalingv2.MetricTarget{Type: autoscalingv2.ValueMetricType, Value: &value}, nil
}

if target.AverageValue != nil {
averageValue := target.AverageValue.DeepCopy()
return autoscalingv2.MetricTarget{Type: autoscalingv2.AverageValueMetricType, AverageValue: &averageValue}, nil
}

return autoscalingv2.MetricTarget{
Type: autoscalingv2.UtilizationMetricType,
AverageUtilization: new(*target.AverageUtilization),
}, nil
}

// SetupWithManager sets up the controller with the Manager.
func (r *WorkloadDeploymentHPAReconciler) SetupWithManager(mgr mcmanager.Manager) error {
r.mgr = mgr

return mcbuilder.ControllerManagedBy(mgr).
Named("workload-deployment-hpa").
For(&computev1alpha.WorkloadDeployment{}, mcbuilder.WithEngageWithLocalCluster(false)).
Owns(&autoscalingv2.HorizontalPodAutoscaler{}, mcbuilder.WithPredicates(predicate.GenerationChangedPredicate{})).
Complete(r)
}
Loading
Loading