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
26 changes: 26 additions & 0 deletions internal/controller/controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,8 @@ import (
kubernetesmetrics "k8s.io/metrics/pkg/client/clientset/versioned"
)

const scaleActivityWindow = 5 * time.Second

type namespaceLister struct {
podIndexer cache.Indexer
podLister listerv1.PodLister
Expand All @@ -51,6 +53,7 @@ type Controller struct {
scaleMu *xsync.Map[function.Function, *sync.Mutex]
routerHeartbeats *xsync.Map[function.Function, RouterHeartbeats]
stabilizationWindows *xsync.Map[function.Function, *StabilizationWindow]
scaleActivity *xsync.Map[function.Function, time.Time]
}

func New(newClientFunc NewClientFunc, kubernetes kubernetes.Interface, kubernetesMetrics kubernetesmetrics.Interface) *Controller {
Expand All @@ -64,6 +67,7 @@ func New(newClientFunc NewClientFunc, kubernetes kubernetes.Interface, kubernete
scaleMu: xsync.NewMap[function.Function, *sync.Mutex](),
routerHeartbeats: xsync.NewMap[function.Function, RouterHeartbeats](),
stabilizationWindows: xsync.NewMap[function.Function, *StabilizationWindow](),
scaleActivity: xsync.NewMap[function.Function, time.Time](),
}
}

Expand Down Expand Up @@ -99,6 +103,28 @@ func (ctrl *Controller) getControllerClient(ip string) Client {
return controllerClient
}

func (ctrl *Controller) markScaleActivity(fn function.Function) {
if ctrl.scaleActivity == nil {
return
}
ctrl.scaleActivity.Store(fn, time.Now())
}

func (ctrl *Controller) lastScaleActivity(fn function.Function) (time.Time, bool) {
if ctrl.scaleActivity == nil {
return time.Time{}, false
}
return ctrl.scaleActivity.Load(fn)
}

func (ctrl *Controller) isRecentlyScaling(fn function.Function) bool {
last, ok := ctrl.lastScaleActivity(fn)
if !ok {
return false
}
return time.Since(last) <= scaleActivityWindow
}

func (ctrl *Controller) startInformers(ctx context.Context) error {
ctx, span := telemetry.Trace(ctx, "controller.start_informers")
defer span.End()
Expand Down
26 changes: 17 additions & 9 deletions internal/controller/scale.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,28 +46,28 @@ var (
Subsystem: "controller",
Name: "waiting_for_unassigned_pods",
Help: "The number of functions that are waiting for an unassigned pod",
}, []string{"function_deployment"})
}, []string{"function_deployment", "function_tenant"})

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

How do feel about moving these function_tenant label additions to a separate PR?

I avoided adding the function_tenant label in the first place because it's technically unbounded, we can have an infinite number of tenants, and prometheus advises against using labels with unbounded values:

Image


assignmentsTotal = promauto.NewCounterVec(prometheus.CounterOpts{
Namespace: "skipper",
Subsystem: "controller",
Name: "assignments_total",
Help: "The number of times the controller has assigned a pod to a function",
}, []string{"function_deployment"})
}, []string{"function_deployment", "function_tenant"})

scaleUpsTotal = promauto.NewCounterVec(prometheus.CounterOpts{
Namespace: "skipper",
Subsystem: "controller",
Name: "scale_ups_total",
Help: "The number of times the controller has scaled up a function",
}, []string{"function_deployment"})
}, []string{"function_deployment", "function_tenant"})

scaleDownsTotal = promauto.NewCounterVec(prometheus.CounterOpts{
Namespace: "skipper",
Subsystem: "controller",
Name: "scale_downs_total",
Help: "The number of times the controller has scaled down a function",
}, []string{"function_deployment"})
}, []string{"function_deployment", "function_tenant"})
)

func (ctrl *Controller) scaleNamespace(ctx context.Context, namespace string) error {
Expand Down Expand Up @@ -317,7 +317,8 @@ func (ctrl *Controller) scale(ctx context.Context, fn function.Function, decisio
}

log.Info(ctx, "scaling function up")
scaleUpsTotal.WithLabelValues(fn.Deployment).Add(float64(decision.DesiredInstances - len(readyInstances)))
ctrl.markScaleActivity(fn)
scaleUpsTotal.WithLabelValues(fn.Deployment, fn.Tenant).Add(float64(decision.DesiredInstances - len(readyInstances)))

for range decision.DesiredInstances - len(readyInstances) {
instance, err := ctrl.assignPod(ctx, fn)
Expand All @@ -329,7 +330,7 @@ func (ctrl *Controller) scale(ctx context.Context, fn function.Function, decisio
} else {
// we either need to scale down or we're already at the desired number of instances but have extra unready instances
log.Info(ctx, "scaling function down")
scaleDownsTotal.WithLabelValues(fn.Deployment).Add(float64(len(readyInstances) + len(unreadyInstances) - decision.DesiredInstances))
scaleDownsTotal.WithLabelValues(fn.Deployment, fn.Tenant).Add(float64(len(readyInstances) + len(unreadyInstances) - decision.DesiredInstances))

// delete all unready instances
for _, unreadyInstance := range unreadyInstances {
Expand Down Expand Up @@ -360,7 +361,7 @@ func (ctrl *Controller) assignPod(ctx context.Context, fn function.Function) (in
ctx, span := telemetry.Trace(ctx, "controller.assign_pod")
defer span.End()

assignmentsTotal.WithLabelValues(fn.Deployment).Inc()
assignmentsTotal.WithLabelValues(fn.Deployment, fn.Tenant).Inc()

GET_UNASSIGNED_POD:
var pod *v1.Pod
Expand Down Expand Up @@ -469,8 +470,8 @@ func (ctrl *Controller) getUnassignedPod(ctx context.Context, fn function.Functi
ctx, span := telemetry.Trace(ctx, "controller.get_unassigned_pod")
defer span.End()

waitingForUnassignedPods.WithLabelValues(fn.Deployment).Inc()
defer waitingForUnassignedPods.WithLabelValues(fn.Deployment).Dec()
waitingForUnassignedPods.WithLabelValues(fn.Deployment, fn.Tenant).Inc()
defer waitingForUnassignedPods.WithLabelValues(fn.Deployment, fn.Tenant).Dec()

return timer.Poll(ctx, 250*time.Millisecond, func(ctx context.Context) (*v1.Pod, error) {
unassignedPods, err := ctrl.getUnassignedPods(fn)
Expand Down Expand Up @@ -745,12 +746,19 @@ type RouterHeartbeats map[string]function.Heartbeat
// Combined returns a heartbeat that is the sum of all the heartbeats from all the routers
func (r RouterHeartbeats) Combined() function.Heartbeat {
var combined function.Heartbeat
combined.InFlightPerInstance = make(map[string]int)
for _, heartbeat := range r {
combined.Function = heartbeat.Function
combined.InFlightRequests += heartbeat.InFlightRequests
if combined.Timestamp.Before(heartbeat.Timestamp) {
combined.Timestamp = heartbeat.Timestamp
}
for instance, count := range heartbeat.InFlightPerInstance {
combined.InFlightPerInstance[instance] += count
}
}
if len(combined.InFlightPerInstance) == 0 {
combined.InFlightPerInstance = map[string]int{}
}
Comment on lines +760 to 762

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: This is redundant.

Suggested change
if len(combined.InFlightPerInstance) == 0 {
combined.InFlightPerInstance = map[string]int{}
}

return combined
}
203 changes: 180 additions & 23 deletions internal/controller/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package controller

import (
"context"
"math"
"math/rand"
"net/http"
"slices"
Expand All @@ -19,12 +20,19 @@ import (
"go.opentelemetry.io/otel/attribute"
)

const (

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These numbers were all picked at random

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

instanceOverloadHeadroom = 0.25 // 25% headroom for instance overload
instanceThrottleInterval = 100 * time.Millisecond
instanceThrottleMaxWait = 2 * time.Second
instanceThrottleDeadlineBuffer = 250 * time.Millisecond
)

var heartbeatsCounter = promauto.NewCounterVec(prometheus.CounterOpts{
Namespace: "skipper",
Subsystem: "controller",
Name: "heartbeats_total",
Help: "The number of heartbeats received by the controller",
}, []string{"function_deployment"})
}, []string{"function_deployment", "function_tenant"})

func (ctrl *Controller) Handler() http.Handler {
mux := http.NewServeMux()
Expand Down Expand Up @@ -54,35 +62,63 @@ func (ctrl *Controller) handleInstance(rw http.ResponseWriter, req *http.Request
ctx = log.With(ctx, key.Function.Field(fn))
ctx = telemetry.WithPropagatedAttributes(ctx, key.Function.Attributes(fn)...)

instances, err := ctrl.getReadyInstances(fn)
if err != nil {
log.Error(ctx, "failed to get instances", key.Error.Field(err))
http.Error(rw, err.Error(), http.StatusInternalServerError)
return
}

telemetry.SetAttributes(ctx, attribute.Bool("has_instances", len(instances) > 0))
selectionStart := time.Now()
var (
instances []*function.Instance
loads map[string]int
)

for len(instances) == 0 {
if instances, err = ctrl.scale(ctx, fn, ScalingDecision{
DesiredInstances: 1,
UnclampedDesiredInstances: 1,
Reason: "no ready instances",
}); err != nil {
log.Error(ctx, "failed to scale function", key.Error.Field(err))
for {
instances, err = ctrl.getReadyInstances(fn)
if err != nil {
log.Error(ctx, "failed to get instances", key.Error.Field(err))
http.Error(rw, err.Error(), http.StatusInternalServerError)
return
}

telemetry.SetAttributes(ctx, attribute.Bool("has_instances", len(instances) > 0))

if len(instances) == 0 {
if _, err = ctrl.scale(ctx, fn, ScalingDecision{
DesiredInstances: 1,
UnclampedDesiredInstances: 1,
Reason: "no ready instances",
}); err != nil {
log.Error(ctx, "failed to scale function", key.Error.Field(err))
http.Error(rw, err.Error(), http.StatusInternalServerError)
return
}
continue
}

if len(instances) > fn.Scale.MaxInstances {
// sort instances by assigned at in descending order (newest first)
slices.SortFunc(instances, func(a, b *function.Instance) int { return b.AssignedAt.Compare(a.AssignedAt) })
// keep the newest instances
instances = instances[:fn.Scale.MaxInstances]
}

loads = ctrl.inFlightPerInstance(fn)

if ctrl.shouldThrottle(ctx, fn, instances, loads, selectionStart) {
totalLoad := totalInFlight(loads)
log.Debug(ctx, "delaying instance selection while scaling", key.InFlightRequests.Field(totalLoad), key.Count.Field(len(instances)))
if !ctrl.waitForCapacity(ctx, selectionStart) {
log.Debug(ctx, "unable to wait longer for additional capacity", key.InFlightRequests.Field(totalLoad))
break
}
continue
}

break
}

if len(instances) > fn.Scale.MaxInstances {
// sort instances by assigned at in descending order (newest first)
slices.SortFunc(instances, func(a, b *function.Instance) int { return b.AssignedAt.Compare(a.AssignedAt) })
// keep the newest instances
instances = instances[:fn.Scale.MaxInstances]
if len(instances) == 0 {
http.Error(rw, "no ready instances", http.StatusServiceUnavailable)
return
}
Comment on lines +116 to 119

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: Let's log the error before returning it so we can easily track when this happens in Axiom.

Suggested change
if len(instances) == 0 {
http.Error(rw, "no ready instances", http.StatusServiceUnavailable)
return
}
if len(instances) == 0 {
log.Error(ctx, "no ready instances")
http.Error(rw, "no ready instances", http.StatusServiceUnavailable)
return
}


instance := instances[rand.Intn(len(instances))]
instance := ctrl.chooseLeastBusyInstance(fn, instances, loads)

rw.Header().Set("Content-Type", "application/json")
rw.WriteHeader(http.StatusOK)
Expand All @@ -91,6 +127,127 @@ func (ctrl *Controller) handleInstance(rw http.ResponseWriter, req *http.Request
}
}

func (ctrl *Controller) inFlightPerInstance(fn function.Function) map[string]int {
routerHeartbeats, ok := ctrl.routerHeartbeats.Load(fn)
if !ok {
return map[string]int{}
}

combined := routerHeartbeats.Combined()
if len(combined.InFlightPerInstance) == 0 {
return map[string]int{}
}

loads := make(map[string]int, len(combined.InFlightPerInstance))
for instance, count := range combined.InFlightPerInstance {
loads[instance] = count
}
return loads
Comment on lines +136 to +145

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: simpler to just do this

Suggested change
combined := routerHeartbeats.Combined()
if len(combined.InFlightPerInstance) == 0 {
return map[string]int{}
}
loads := make(map[string]int, len(combined.InFlightPerInstance))
for instance, count := range combined.InFlightPerInstance {
loads[instance] = count
}
return loads
return maps.Clone(routerHeartbeats.Combined().InFlightPerInstance)

}

func (ctrl *Controller) shouldThrottle(ctx context.Context, fn function.Function, instances []*function.Instance, loads map[string]int, start time.Time) bool {
if fn.Scale.TargetInFlightRequests <= 0 {
return false
}
if len(instances) == 0 {
return false
}

threshold := int(math.Ceil(float64(fn.Scale.TargetInFlightRequests) * (1 + instanceOverloadHeadroom)))
if threshold < 0 {
return false
}

for _, instance := range instances {
if loads[instance.Name] <= threshold {
return false
}
}

if !ctrl.isRecentlyScaling(fn) {
return false
}
Comment on lines +167 to +169

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I would remove this ctrl.isRecentlyScaling check and all its dependencies imo.

  1. ctrl.markScaleActivity is only called on the controller that is responsible for the function in the hashring, so this could return false when it's really true.
  2. We know all the existing instances are above their threshold and giving them a moment to complete some requests should help them recover regardless of whether we just scaled or not


return ctrl.canThrottle(ctx, start)
}

func (ctrl *Controller) canThrottle(ctx context.Context, start time.Time) bool {
if time.Since(start) >= instanceThrottleMaxWait {
return false
}

if deadline, ok := ctx.Deadline(); ok {
if time.Until(deadline) <= instanceThrottleDeadlineBuffer {
return false
}
}

return true
}

func (ctrl *Controller) waitForCapacity(ctx context.Context, start time.Time) bool {
elapsed := time.Since(start)
if elapsed >= instanceThrottleMaxWait {
return false
}

wait := instanceThrottleInterval
if remaining := instanceThrottleMaxWait - elapsed; remaining < wait {
wait = remaining
}
if wait <= 0 {
return false
}

select {
case <-ctx.Done():
return false
case <-time.After(wait):
return true
}
}

func (ctrl *Controller) chooseLeastBusyInstance(fn function.Function, instances []*function.Instance, loads map[string]int) *function.Instance {
if len(instances) == 0 {
return nil
}

var candidates []*function.Instance

// If no threshold is set, use all instances
if fn.Scale.TargetInFlightRequests <= 0 {
candidates = instances
} else {
// Calculate threshold: instances above this are considered overloaded
threshold := int(math.Ceil(float64(fn.Scale.TargetInFlightRequests) * (1 + instanceOverloadHeadroom)))

// Filter out instances above the threshold
candidates = make([]*function.Instance, 0, len(instances))
for _, instance := range instances {
load := loads[instance.Name]
if load <= threshold {
candidates = append(candidates, instance)
}
}

// If no instances are below threshold, fall back to all instances
if len(candidates) == 0 {
candidates = instances
}
}

// Randomly select from the acceptable instances
return candidates[rand.Intn(len(candidates))]
}

func totalInFlight(loads map[string]int) int {
total := 0
for _, count := range loads {
total += count
}
return total
}

func (ctrl *Controller) handleScale(rw http.ResponseWriter, req *http.Request) {
ctx := req.Context()
fn, err := function.FromHeader(req)
Expand Down Expand Up @@ -147,7 +304,7 @@ func (ctrl *Controller) handleHeartbeat(rw http.ResponseWriter, req *http.Reques
}

for _, heartbeat := range heartbeats {
heartbeatsCounter.WithLabelValues(heartbeat.Function.Deployment).Inc()
heartbeatsCounter.WithLabelValues(heartbeat.Function.Deployment, heartbeat.Function.Tenant).Inc()

ctrl.routerHeartbeats.Compute(heartbeat.Function, func(routerHeartbeats RouterHeartbeats, loaded bool) (RouterHeartbeats, xsync.ComputeOp) {
if !loaded {
Expand Down
Loading