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
34 changes: 2 additions & 32 deletions internal/controller/pod_placement_helpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -128,13 +128,13 @@ func (r *PodPlacementReconciler) selectPlacement(ctx context.Context, pod *corev
"provider", ref.Name, "capacityType", tier)
continue // unregistered; NodePool status surfaces this separately
}
if !servesCapacityTier(prov, tier) {
if !prov.Capabilities().ServesCapacityTier(tier) {
metrics.RecordCandidateSkip(ref.Name, tier, "", metrics.SkipCapacityUnsupported)
log.V(1).Info("skipping candidate: provider does not offer the capacity tier",
"provider", ref.Name, "capacityType", tier)
continue
}
if !servesEgress(prov, pool.Spec.Egress) {
if !prov.Capabilities().ServesEgress(pool.Spec.Egress) {
metrics.RecordCandidateSkip(ref.Name, tier, "", metrics.SkipEgressUnsupported)
log.V(1).Info("skipping candidate: provider cannot enforce the pool's egress policy",
"provider", ref.Name, "egressMode", pool.Spec.Egress.ModeOrOpen())
Expand Down Expand Up @@ -232,36 +232,6 @@ func capacityTiers(pool *nebulav1alpha1.NodePool) []nebulav1alpha1.CapacityType
return pool.Spec.CapacityTypes
}

// servesCapacityTier reports whether prov can deliver the candidate's capacity tier. Only Spot
// is ever refused: an OnDemand-only provider (Modal) has no interruptible tier, so placing a
// Spot candidate there would stamp CapacityType=Spot on the Pod, hand it to an adapter that
// drops the field, and bill OnDemand rates for capacity the user asked to be cheap — with no
// error or event revealing the substitution. Skipping lets the pool's next tier take over
// ([Spot, OnDemand] still lands on Modal, now truthfully labelled), and a Spot-only pool
// leaves the Pod visibly unplaceable rather than quietly overcharged.
//
// The empty tier is "the provider's default", which every provider serves, so it passes.
func servesCapacityTier(prov provider.Provider, tier nebulav1alpha1.CapacityType) bool {
if tier != nebulav1alpha1.CapacitySpot {
return true
}
return prov.Capabilities().SupportsSpot
}

// servesEgress reports whether prov can enforce the pool's egress policy. Open needs no
// enforcement, so every provider serves it; anything else needs SupportsEgressPolicy.
//
// Same reasoning as servesCapacity, and load-bearing for a different reason: a provider
// that drops the field would put the workload on the open internet while the pool claims
// containment. Skipping makes that visible — an AWS-only pool asking for Blocked leaves the
// Pod unplaceable instead of silently unprotected.
func servesEgress(prov provider.Provider, policy *nebulav1alpha1.EgressPolicy) bool {
if !policy.RestrictsEgress() {
return true
}
return prov.Capabilities().SupportsEgressPolicy
}

// requestedGeographies reads the Pod's RegionsAnnotation into the narrowing ResolveRegions
// takes. Absent means no narrowing, the common case.
func requestedGeographies(pod *corev1.Pod) (narrowTo []string, ok bool) {
Expand Down
82 changes: 54 additions & 28 deletions internal/webhook/v1/nodepool_webhook.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,19 +19,16 @@ package v1
import (
"context"
"fmt"
"strings"

"k8s.io/apimachinery/pkg/runtime"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/webhook"
"sigs.k8s.io/controller-runtime/pkg/webhook/admission"

nebulav1alpha1 "github.com/InftyAI/Nebula/api/v1alpha1"
"github.com/InftyAI/Nebula/pkg/provider"
)

// modalProvider is the provider name Modal registers its virtual node under.
const modalProvider = "modal"

// SetupNodePoolWebhookWithManager registers the webhook for NodePool in the manager.
func SetupNodePoolWebhookWithManager(mgr ctrl.Manager) error {
return ctrl.NewWebhookManagedBy(mgr).For(&nebulav1alpha1.NodePool{}).
Expand All @@ -41,55 +38,84 @@ func SetupNodePoolWebhookWithManager(mgr ctrl.Manager) error {

// +kubebuilder:webhook:path=/validate-nebula-inftyai-com-v1alpha1-nodepool,mutating=false,failurePolicy=fail,sideEffects=None,groups=nebula.inftyai.com,resources=nodepools,verbs=create;update,versions=v1alpha1,name=vnodepool-v1alpha1.nebula.inftyai.com,admissionReviewVersions=v1

// NodePoolCustomValidator rejects NodePools that ask for capacity a provider
// cannot supply. Modal has no spot instances, so a pool that lists Modal and
// Spot would hold a Spot tier that Modal can never fill.
//
// The API server applies the capacityTypes default ([OnDemand, Spot]) before
// admission, so a Modal pool that leaves capacityTypes empty is rejected too;
// it must list [OnDemand] explicitly.
type NodePoolCustomValidator struct{}
// NodePoolCustomValidator validates NodePools on create and update. A setting a listed
// provider cannot serve (Spot on Modal, a restricted egress on RunPod) is NOT rejected:
// placement skips that provider instead (see provider.Capabilities), so the rest of the pool
// still works. It only WARNS, so a provider placement will never use is not a surprise.
type NodePoolCustomValidator struct {
// Providers resolves a provider name to its backend; defaults to the registry. An
// unresolved name draws no warning: placement skips it too, and NodePool status reports it.
Providers func(name string) (provider.Provider, bool)
}

var _ webhook.CustomValidator = &NodePoolCustomValidator{}

// ValidateCreate implements webhook.CustomValidator.
func (v *NodePoolCustomValidator) ValidateCreate(_ context.Context, obj runtime.Object) (admission.Warnings, error) {
return nil, validateNodePool(obj)
return v.validate(obj)
Comment thread
kerthcet marked this conversation as resolved.
}

// ValidateUpdate implements webhook.CustomValidator.
func (v *NodePoolCustomValidator) ValidateUpdate(_ context.Context, _, newObj runtime.Object) (admission.Warnings, error) {
return nil, validateNodePool(newObj)
return v.validate(newObj)
}

// ValidateDelete implements webhook.CustomValidator. Deletes are always allowed.
func (v *NodePoolCustomValidator) ValidateDelete(_ context.Context, _ runtime.Object) (admission.Warnings, error) {
return nil, nil
}

func validateNodePool(obj runtime.Object) error {
func (v *NodePoolCustomValidator) validate(obj runtime.Object) (admission.Warnings, error) {
np, ok := obj.(*nebulav1alpha1.NodePool)
if !ok {
return fmt.Errorf("expected a NodePool object but got %T", obj)
}
if hasProvider(np, modalProvider) && hasCapacityType(np, nebulav1alpha1.CapacitySpot) {
return fmt.Errorf("modal does not support spot instances; remove modal from providers or set capacityTypes to [OnDemand]")
return nil, fmt.Errorf("expected a NodePool object but got %T", obj)
}
return nil
return v.warnings(np), nil
}

func hasProvider(np *nebulav1alpha1.NodePool, name string) bool {
for _, p := range np.Spec.Providers {
if strings.EqualFold(p.Name, name) {
return true
// warnings names every listed provider placement will never use, and says so once more when
// that is all of them, since such a pool admits Pods it can never place.
func (v *NodePoolCustomValidator) warnings(np *nebulav1alpha1.NodePool) admission.Warnings {
var warns admission.Warnings
known, usable := 0, 0
for _, ref := range np.Spec.Providers {
prov, ok := v.provider(ref.Name)
if !ok {
continue
}
known++
caps := prov.Capabilities()
switch {
case !caps.ServesEgress(np.Spec.Egress):
warns = append(warns, fmt.Sprintf("provider %q cannot enforce egress mode %q, so placement will never use it",
ref.Name, np.Spec.Egress.ModeOrOpen()))
case !servesAnyTier(caps, np.Spec.CapacityTypes):
warns = append(warns, fmt.Sprintf("provider %q serves none of capacityTypes %v, so placement will never use it",
ref.Name, np.Spec.CapacityTypes))
default:
usable++
}
}
return false
if known > 0 && usable == 0 {
warns = append(warns, "no provider in this pool can serve it; its Pods will stay unplaceable")
}
return warns
}

func hasCapacityType(np *nebulav1alpha1.NodePool, ct nebulav1alpha1.CapacityType) bool {
for _, c := range np.Spec.CapacityTypes {
if c == ct {
func (v *NodePoolCustomValidator) provider(name string) (provider.Provider, bool) {
if v.Providers != nil {
return v.Providers(name)
}
return provider.Get(name)
}

// servesAnyTier mirrors placement's tier walk, where an empty list is one default tier.
func servesAnyTier(caps provider.Capabilities, tiers []nebulav1alpha1.CapacityType) bool {
if len(tiers) == 0 {
return true
}
for _, t := range tiers {
if caps.ServesCapacityTier(t) {
return true
}
}
Expand Down
86 changes: 62 additions & 24 deletions internal/webhook/v1/nodepool_webhook_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,58 +18,96 @@ package v1

import (
"context"
"strings"
"testing"

metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"

nebulav1alpha1 "github.com/InftyAI/Nebula/api/v1alpha1"
"github.com/InftyAI/Nebula/pkg/provider"
)

func nodePool(capacity []nebulav1alpha1.CapacityType, providers ...string) *nebulav1alpha1.NodePool {
// capsOnly is a provider whose only behaviour is its Capabilities; the webhook reads nothing else.
type capsOnly struct {
provider.Provider
caps provider.Capabilities
}

func (p capsOnly) Capabilities() provider.Capabilities { return p.caps }

// Shaped like the real adapters: Modal enforces egress but has no Spot, RunPod neither, AWS Spot only.
var testProviders = map[string]provider.Provider{
"modal": capsOnly{caps: provider.Capabilities{SupportsEgressPolicy: true}},
"runpod": capsOnly{},
"aws": capsOnly{caps: provider.Capabilities{SupportsSpot: true}},
}

func lookup(name string) (provider.Provider, bool) {
p, ok := testProviders[name]
return p, ok
}

func pool(capacity []nebulav1alpha1.CapacityType, egress nebulav1alpha1.EgressMode, providers ...string) *nebulav1alpha1.NodePool {
np := &nebulav1alpha1.NodePool{ObjectMeta: metav1.ObjectMeta{Name: "np"}}
for _, p := range providers {
np.Spec.Providers = append(np.Spec.Providers, nebulav1alpha1.ProviderSpec{Name: p})
}
np.Spec.CapacityTypes = capacity
if egress != "" {
np.Spec.Egress = &nebulav1alpha1.EgressPolicy{Mode: egress}
}
return np
}

func TestNodePoolValidator(t *testing.T) {
// A setting a provider cannot serve is never an error, only a warning naming the provider
// placement will skip, plus one more when nothing in the pool is left.
func TestNodePoolValidatorWarnings(t *testing.T) {
spot, onDemand := nebulav1alpha1.CapacitySpot, nebulav1alpha1.CapacityOnDemand
both := []nebulav1alpha1.CapacityType{onDemand, spot}
cases := []struct {
name string
np *nebulav1alpha1.NodePool
wantErr bool
name string
np *nebulav1alpha1.NodePool
want []string // substrings, one per expected warning, in order
}{
{"modal with spot", nodePool([]nebulav1alpha1.CapacityType{spot}, "modal"), true},
{"modal among others with spot", nodePool([]nebulav1alpha1.CapacityType{onDemand, spot}, "runpod", "modal"), true},
{"modal name is case-insensitive", nodePool([]nebulav1alpha1.CapacityType{spot}, "Modal"), true},
{"modal with on-demand only", nodePool([]nebulav1alpha1.CapacityType{onDemand}, "modal"), false},
{"other provider with spot", nodePool([]nebulav1alpha1.CapacityType{spot, onDemand}, "runpod"), false},
{"every provider usable", pool(both, "", "modal", "runpod", "aws"), nil},
// Modal still serves the OnDemand tier, so the default capacityTypes must stay quiet.
{"modal with spot as a fallback", pool(both, "", "modal"), nil},
{"modal with spot only", pool([]nebulav1alpha1.CapacityType{spot}, "", "modal", "aws"),
[]string{`"modal" serves none of capacityTypes`}},
{"runpod under a restricted egress", pool(both, nebulav1alpha1.EgressBlocked, "runpod", "modal"),
[]string{`"runpod" cannot enforce egress mode "Blocked"`}},
{"nothing can place", pool(both, nebulav1alpha1.EgressAllowlist, "runpod", "aws"),
[]string{`"runpod" cannot enforce`, `"aws" cannot enforce`, "will stay unplaceable"}},
// Unregistered is placement's skip and status's report, not a warning.
{"unregistered provider", pool(both, "", "lambda"), nil},
{"empty capacityTypes is the default tier", pool(nil, "", "modal"), nil},
}
v := &NodePoolCustomValidator{}
v := &NodePoolCustomValidator{Providers: lookup}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if _, err := v.ValidateCreate(context.Background(), tc.np); (err != nil) != tc.wantErr {
t.Fatalf("ValidateCreate err = %v, wantErr %v", err, tc.wantErr)
warns, err := v.ValidateCreate(context.Background(), tc.np)
if err != nil {
t.Fatalf("ValidateCreate err = %v, want nil", err)
}
if len(warns) != len(tc.want) {
t.Fatalf("warnings = %q, want %d matching %q", warns, len(tc.want), tc.want)
}
for i, w := range tc.want {
if !strings.Contains(warns[i], w) {
t.Errorf("warning %d = %q, want it to contain %q", i, warns[i], w)
}
}
old := nodePool([]nebulav1alpha1.CapacityType{onDemand}, "modal")
if _, err := v.ValidateUpdate(context.Background(), old, tc.np); (err != nil) != tc.wantErr {
t.Fatalf("ValidateUpdate err = %v, wantErr %v", err, tc.wantErr)
upd, err := v.ValidateUpdate(context.Background(), pool(both, "", "modal"), tc.np)
if err != nil || len(upd) != len(warns) {
t.Errorf("ValidateUpdate = %q, %v; want the create result %q", upd, err, warns)
}
})
}
}

func TestNodePoolValidatorAllowsDelete(t *testing.T) {
np := nodePool([]nebulav1alpha1.CapacityType{nebulav1alpha1.CapacitySpot}, "modal")
if _, err := (&NodePoolCustomValidator{}).ValidateDelete(context.Background(), np); err != nil {
t.Fatalf("ValidateDelete err = %v, want nil", err)
}
}

func TestNodePoolValidatorRejectsWrongType(t *testing.T) {
if _, err := (&NodePoolCustomValidator{}).ValidateCreate(context.Background(), &nebulav1alpha1.NodePoolList{}); err == nil {
v := &NodePoolCustomValidator{Providers: lookup}
if _, err := v.ValidateCreate(context.Background(), &nebulav1alpha1.NodePoolList{}); err == nil {
t.Fatal("expected an error for a non-NodePool object")
}
}
4 changes: 2 additions & 2 deletions pkg/provider/catalog/base.go
Original file line number Diff line number Diff line change
Expand Up @@ -149,8 +149,8 @@ func (b Base) PricePerHour(req provider.PriceRequest) (float64, error) {
return 0, fmt.Errorf("catalog: %s: accelerator %q with count %d", b.ProviderName, req.AcceleratorType, req.Count)
}
// An empty tier is the candidate a pool declaring no capacityTypes emits.
// servesCapacityTier already reads that as non-Spot, so resolve it here rather than let
// row order decide — on AWS that is a threefold difference.
// Capabilities.ServesCapacityTier already reads that as non-Spot, so resolve it here
// rather than let row order decide — on AWS that is a threefold difference.
tier := req.CapacityType
if tier == "" {
tier = nebulav1alpha1.CapacityOnDemand
Expand Down
24 changes: 24 additions & 0 deletions pkg/provider/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -340,6 +340,30 @@ type Capabilities struct {
ProvisionTimeout time.Duration
}

// ServesCapacityTier reports whether the provider can deliver tier. Only Spot is ever refused:
// an OnDemand-only provider (Modal) has no interruptible tier, so placing a Spot candidate
// there would stamp CapacityType=Spot on the Pod, hand it to an adapter that drops the field,
// and bill OnDemand rates for capacity the user asked to be cheap — with no error or event
// revealing the substitution. Skipping lets the pool's next tier take over ([Spot, OnDemand]
// still lands on Modal, now truthfully labelled), and a Spot-only pool leaves the Pod visibly
// unplaceable rather than quietly overcharged.
//
// The empty tier is "the provider's default", which every provider serves, so it passes.
func (c Capabilities) ServesCapacityTier(tier nebulav1alpha1.CapacityType) bool {
return tier != nebulav1alpha1.CapacitySpot || c.SupportsSpot
}

// ServesEgress reports whether the provider can enforce policy. Open needs no enforcement,
// so every provider serves it; anything else needs SupportsEgressPolicy.
//
// Load-bearing for a different reason than ServesCapacityTier: a provider that drops the
// field would put the workload on the open internet while the pool claims containment.
// Skipping makes that visible — an AWS-only pool asking for Blocked leaves the Pod
// unplaceable instead of silently unprotected.
func (c Capabilities) ServesEgress(policy *nebulav1alpha1.EgressPolicy) bool {
return !policy.RestrictsEgress() || c.SupportsEgressPolicy
}

// Instance is the provider-agnostic view of one external instance, as observed.
type Instance struct {
ID string
Expand Down
Loading