diff --git a/internal/controller/nodeclaim_controller.go b/internal/controller/nodeclaim_controller.go index 336bdad..6f09310 100644 --- a/internal/controller/nodeclaim_controller.go +++ b/internal/controller/nodeclaim_controller.go @@ -411,6 +411,7 @@ func (r *NodeClaimReconciler) recordPrice(ctx context.Context, nc *nebulav1alpha CapacityType: nc.Spec.CapacityType, CPUCores: cpuCores, MemoryMiB: memoryMiB, + DiskGiB: util.PodEphemeralStorageGiB(pod), }) if err != nil { if errors.Is(err, provider.ErrNoPrice) { diff --git a/pkg/provider/aws/aws.go b/pkg/provider/aws/aws.go index 37a5f9a..80f7cce 100644 --- a/pkg/provider/aws/aws.go +++ b/pkg/provider/aws/aws.go @@ -50,6 +50,7 @@ import ( nebulav1alpha1 "github.com/InftyAI/Nebula/api/v1alpha1" "github.com/InftyAI/Nebula/pkg/provider" "github.com/InftyAI/Nebula/pkg/provider/catalog" + "github.com/InftyAI/Nebula/pkg/provider/catalog/data" "github.com/InftyAI/Nebula/pkg/util" ) @@ -167,6 +168,9 @@ type InstanceSpec struct { Region string // Tags carry Nebula identity; ClaimTagKey holds the NodeClaim name. Tags map[string]string + // DiskGiB is the user space added to the root volume's OS base, from + // util.PodEphemeralStorageGiB; see sdkClient.rootVolume. + DiskGiB int } // EC2Instance is the adapter-level view of one EC2 instance as observed. @@ -739,6 +743,29 @@ func (p *Provider) ClassifyProvisionError(err error, accelerator, region string) return scope } +// awsAMIRootGiB is the GPU AMI's root snapshot (30 GiB in every region checked, 2026-10-04), +// the OS base of every root volume. A constant because pricing has no region to resolve the +// AMI in; a larger snapshot is logged at client construction, since it under-prices. +const awsAMIRootGiB = 30 + +// awsMaxDiskGiB is the most user space a Pod may ask for: gp3's 16 TiB volume cap less the OS +// base. Refused at launch and unpriced above it, so it never reaches rootVolume's int32. +const awsMaxDiskGiB = 16*1024 - awsAMIRootGiB + +// PricePerHour overrides catalog.Base to add the root volume, which EBS bills by provisioned +// size apart from the instance. It prices the size sdkClient.rootVolume launches. +func (p *Provider) PricePerHour(req provider.PriceRequest) (float64, error) { + if req.DiskGiB > awsMaxDiskGiB { + return 0, fmt.Errorf("aws: %d GiB disk exceeds the %d GiB a root volume can add: %w", + req.DiskGiB, awsMaxDiskGiB, provider.ErrNoPrice) + } + rate, err := p.Base.PricePerHour(req) + if err != nil { + return 0, err + } + return rate + data.AWSRootVolumeCostPerHour(awsAMIRootGiB+req.DiskGiB), nil +} + // instanceSpecFromPod reads the workload off the Pod (source of truth) and the // accelerator type (from the AcceleratorTypeLabel), maps it to an EC2 instance // type via the catalog, and stamps the claim tag, capacity tier, and region. @@ -778,6 +805,10 @@ func (p *Provider) instanceSpecFromPod( return InstanceSpec{}, errors.New( "aws: pod requests no accelerator; EC2 GPU provisioning needs an accelerator type and count") } + diskGiB := util.PodEphemeralStorageGiB(pod) + if diskGiB > awsMaxDiskGiB { + return InstanceSpec{}, fmt.Errorf("aws: %d GiB disk exceeds the %d GiB a root volume can add", diskGiB, awsMaxDiskGiB) + } instanceTypes, ok := p.MapAccelerator(canonical, count) if !ok { return InstanceSpec{}, fmt.Errorf("aws: no EC2 instance type for %s x%d", canonical, count) @@ -800,10 +831,11 @@ func (p *Provider) instanceSpecFromPod( // TODO: deliver Secret-derived values out-of-band — SSM Parameter Store / Secrets // Manager under the claim, fetched at boot with the instance profile — and keep only // non-sensitive values in user-data. - Env: req.Env, - Spot: req.CapacityType == nebulav1alpha1.CapacitySpot, - Region: req.Region, - Tags: map[string]string{ClaimTagKey: req.ClaimName}, + Env: req.Env, + Spot: req.CapacityType == nebulav1alpha1.CapacitySpot, + Region: req.Region, + Tags: map[string]string{ClaimTagKey: req.ClaimName}, + DiskGiB: diskGiB, }, nil } diff --git a/pkg/provider/aws/aws_test.go b/pkg/provider/aws/aws_test.go index 593544c..4137678 100644 --- a/pkg/provider/aws/aws_test.go +++ b/pkg/provider/aws/aws_test.go @@ -36,6 +36,7 @@ import ( nebulav1alpha1 "github.com/InftyAI/Nebula/api/v1alpha1" "github.com/InftyAI/Nebula/pkg/provider" + "github.com/InftyAI/Nebula/pkg/provider/catalog/data" "github.com/InftyAI/Nebula/pkg/util" ) @@ -1037,8 +1038,8 @@ func TestResolveGPUAMI_PicksNewestAndErrsWhenAbsent(t *testing.T) { if err != nil { t.Fatalf("resolveGPUAMI: %v", err) } - if got != "ami-new" { - t.Fatalf("resolveGPUAMI = %q, want ami-new (newest)", got) + if id := awssdk.ToString(got.ImageId); id != "ami-new" { + t.Fatalf("resolveGPUAMI = %q, want ami-new (newest)", id) } // No matching image => ErrConfig (AWS unusable in the region, non-fatal skip). @@ -1070,3 +1071,62 @@ func TestDiscoverDefaultSubnets_ReturnsPerAZTargets(t *testing.T) { t.Fatalf("discoverDefaultSubnets(no default VPC) = (%+v, %v), want (nil, nil)", got, err) } } + +func TestPricePerHour_AddsRootVolume(t *testing.T) { + p := newTestProvider(&fakeClient{}) + req := provider.PriceRequest{AcceleratorType: "T4", Count: 1, CapacityType: nebulav1alpha1.CapacityOnDemand} + + // The OS base is billed even with no disk requested: EBS charges the provisioned size. + got, err := p.PricePerHour(req) + if want := 0.526 + data.AWSRootVolumeCostPerHour(awsAMIRootGiB); err != nil || got != want { + t.Fatalf("PricePerHour(no disk) = %v, %v; want instance + OS base %v", got, err, want) + } + req.DiskGiB = 100 + got, err = p.PricePerHour(req) + if want := 0.526 + data.AWSRootVolumeCostPerHour(awsAMIRootGiB+100); err != nil || got != want { + t.Fatalf("PricePerHour(100 GiB) = %v, %v; want instance + base + 100 GiB %v", got, err, want) + } + // No instance price means no price at all, not a disk-only one. + _, err = p.PricePerHour(provider.PriceRequest{AcceleratorType: "B200", Count: 8, DiskGiB: 100}) + if !errors.Is(err, provider.ErrNoPrice) { + t.Fatalf("PricePerHour(unknown accelerator) err = %v, want ErrNoPrice", err) + } + req.DiskGiB = awsMaxDiskGiB + 1 + if _, err = p.PricePerHour(req); !errors.Is(err, provider.ErrNoPrice) { + t.Fatalf("PricePerHour(above awsMaxDiskGiB) err = %v, want ErrNoPrice", err) + } +} + +func TestProvision_SizesDiskFromEphemeralStorage(t *testing.T) { + f := &fakeClient{runID: "i-disk"} + p := newTestProvider(f) + pod := gpuPod("T4", 1) + pod.Spec.Containers[0].Resources.Limits[corev1.ResourceEphemeralStorage] = resource.MustParse("200Gi") + + if _, err := p.Provision(context.Background(), pod, provider.ProvisionRequest{ + ClaimName: "claim-disk", + Region: "us-west-2", + }); err != nil { + t.Fatalf("Provision: %v", err) + } + if got := f.lastSpec.DiskGiB; got != 200 { + t.Fatalf("spec DiskGiB = %d, want 200", got) + } +} + +func TestProvision_RefusesDiskAboveVolumeCap(t *testing.T) { + f := &fakeClient{runID: "i-disk"} + p := newTestProvider(f) + pod := gpuPod("T4", 1) + pod.Spec.Containers[0].Resources.Limits[corev1.ResourceEphemeralStorage] = resource.MustParse("16Ti") + + if _, err := p.Provision(context.Background(), pod, provider.ProvisionRequest{ + ClaimName: "claim-disk", + Region: "us-west-2", + }); err == nil { + t.Fatal("Provision of 16 TiB user space succeeded, want refusal: the OS base pushes it past gp3's cap") + } + if f.runCnt != 0 { + t.Fatalf("RunInstance called %d times for a refused disk", f.runCnt) + } +} diff --git a/pkg/provider/aws/client.go b/pkg/provider/aws/client.go index ec4c6a7..0f51e7f 100644 --- a/pkg/provider/aws/client.go +++ b/pkg/provider/aws/client.go @@ -119,6 +119,10 @@ type sdkClient struct { // amiID is the region's GPU AMI, resolved at construction. Every instance // launches from it; it is NON-SECRET, AWS-published config, not a credential. amiID string + // rootDevice and rootGiB are the AMI's root device name and snapshot size, which a + // resized root volume must reuse and may not shrink below (see rootVolume). + rootDevice string + rootGiB int // subnets are the default VPC's per-AZ subnets RunInstance fails over across on // a capacity error, discovered at construction. Empty when the region has no // default VPC: RunInstance then makes a single attempt letting EC2 pick the @@ -210,11 +214,16 @@ func newSDKClientForRegion(ctx context.Context, region string) (Client, error) { // Resolve the region's GPU AMI (required — no AMI, nothing to launch) and the // default VPC's per-AZ subnets (best-effort — no default VPC leaves c.subnets // empty and RunInstance lets EC2 pick the subnet, just without zone failover). - amiID, err := c.resolveGPUAMI(ctx) + ami, err := c.resolveGPUAMI(ctx) if err != nil { return nil, fmt.Errorf("aws: resolve GPU AMI in %s: %w", cfg.Region, err) } - c.amiID = amiID + c.amiID = *ami.ImageId + c.rootDevice, c.rootGiB = rootDeviceOf(ami) + if c.rootGiB > awsAMIRootGiB { + logf.FromContext(ctx).Info("AMI root exceeds awsAMIRootGiB; instances are under-priced by the difference", + "region", cfg.Region, "ami", c.amiID, "rootGiB", c.rootGiB, "awsAMIRootGiB", awsAMIRootGiB) + } subnets, err := c.discoverDefaultSubnets(ctx) if err != nil { @@ -436,8 +445,9 @@ func (c *sdkClient) createLaunchTemplate(ctx context.Context, spec InstanceSpec, // template tier-agnostic, so the stable per-claim template is safely reused when // failover retries the same claim under the other tier. ltData := &ec2types.RequestLaunchTemplateData{ - ImageId: awssdk.String(c.amiID), - UserData: awssdk.String(userData), + ImageId: awssdk.String(c.amiID), + UserData: awssdk.String(userData), + BlockDeviceMappings: c.rootVolume(spec.DiskGiB), TagSpecifications: []ec2types.LaunchTemplateTagSpecificationRequest{{ ResourceType: ec2types.ResourceTypeInstance, Tags: ec2Tags(spec.Tags), @@ -620,7 +630,7 @@ func fleetErrorRank(code string) int { // the self-configuring model. A region that returns no matching image is a config // error (ErrConfig): AWS is effectively not usable there, and the caller skips it // non-fatally rather than launching from a missing AMI. -func (c *sdkClient) resolveGPUAMI(ctx context.Context) (string, error) { +func (c *sdkClient) resolveGPUAMI(ctx context.Context) (ec2types.Image, error) { out, err := c.ec2.DescribeImages(ctx, &ec2.DescribeImagesInput{ Owners: []string{"amazon"}, Filters: []ec2types.Filter{ @@ -629,7 +639,7 @@ func (c *sdkClient) resolveGPUAMI(ctx context.Context) (string, error) { }, }) if err != nil { - return "", err + return ec2types.Image{}, err } // Pick the newest by CreationDate (RFC3339 strings sort lexicographically in // chronological order), so a driver/runtime refresh is picked up automatically. @@ -643,9 +653,36 @@ func (c *sdkClient) resolveGPUAMI(ctx context.Context) (string, error) { } } if newest.ImageId == nil { - return "", fmt.Errorf("no GPU AMI (%s) offered: %w", gpuAMINameFilter, ErrConfig) + return ec2types.Image{}, fmt.Errorf("no GPU AMI (%s) offered: %w", gpuAMINameFilter, ErrConfig) } - return *newest.ImageId, nil + return newest, nil +} + +// rootDeviceOf returns the image's root device name and its snapshot size in GiB, zero +// values when the image does not describe one. +func rootDeviceOf(img ec2types.Image) (device string, sizeGiB int) { + device = awssdk.ToString(img.RootDeviceName) + for _, m := range img.BlockDeviceMappings { + if awssdk.ToString(m.DeviceName) == device && m.Ebs != nil { + return device, int(awssdk.ToInt32(m.Ebs.VolumeSize)) + } + } + return device, 0 +} + +// rootVolume is the launch template's root volume: the OS base (awsAMIRootGiB, or the AMI's +// snapshot if larger, which EC2 requires) plus diskGiB of user space for the pulled image and +// the workload's writes. Always sent, even for 0, so the type is gp3, the one PricePerHour +// charges for; the AMI's own root is gp2. +func (c *sdkClient) rootVolume(diskGiB int) []ec2types.LaunchTemplateBlockDeviceMappingRequest { + return []ec2types.LaunchTemplateBlockDeviceMappingRequest{{ + DeviceName: awssdk.String(c.rootDevice), + Ebs: &ec2types.LaunchTemplateEbsBlockDeviceRequest{ + VolumeSize: awssdk.Int32(int32(max(awsAMIRootGiB, c.rootGiB) + diskGiB)), + VolumeType: ec2types.VolumeTypeGp3, + DeleteOnTermination: awssdk.Bool(true), + }, + }} } // discoverDefaultSubnets lists the default VPC's subnets — one default subnet per diff --git a/pkg/provider/aws/client_test.go b/pkg/provider/aws/client_test.go index 6e9b6c4..ac340a6 100644 --- a/pkg/provider/aws/client_test.go +++ b/pkg/provider/aws/client_test.go @@ -1262,3 +1262,52 @@ func TestSDKList_StatusProbeFailureIsNonFatal(t *testing.T) { t.Fatalf("list = %+v, want the instance returned with checks not passed", list) } } + +func TestRootDeviceOf(t *testing.T) { + img := ec2types.Image{ + RootDeviceName: awssdk.String("/dev/xvda"), + BlockDeviceMappings: []ec2types.BlockDeviceMapping{ + {DeviceName: awssdk.String("/dev/sdb"), Ebs: &ec2types.EbsBlockDevice{VolumeSize: awssdk.Int32(500)}}, + {DeviceName: awssdk.String("/dev/xvda"), Ebs: &ec2types.EbsBlockDevice{VolumeSize: awssdk.Int32(30)}}, + }, + } + if dev, size := rootDeviceOf(img); dev != "/dev/xvda" || size != 30 { + t.Fatalf("rootDeviceOf = %q, %d; want /dev/xvda, 30 (the root mapping, not the first)", dev, size) + } + if dev, size := rootDeviceOf(ec2types.Image{}); dev != "" || size != 0 { + t.Fatalf("rootDeviceOf(empty) = %q, %d; want zero values", dev, size) + } +} + +func TestSDKRunInstance_SizesRootVolume(t *testing.T) { + cases := map[string]struct { + disk, amiRoot int + wantSize int32 + }{ + "unset is the OS base alone, still gp3": {0, 30, 30}, + "user space adds to the base": {10, 30, 40}, + "a larger snapshot raises the base": {10, 50, 60}, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + f := &fakeEC2{fleetOut: fleetWith("i-1")} + c := &sdkClient{ec2: f, region: testRegion, amiID: "ami-123", rootDevice: "/dev/xvda", rootGiB: tc.amiRoot} + if _, err := c.RunInstance(context.Background(), InstanceSpec{ + InstanceTypes: []string{"g4dn.xlarge"}, Image: "img", DiskGiB: tc.disk, + Tags: map[string]string{ClaimTagKey: "c"}, + }); err != nil { + t.Fatalf("RunInstance: %v", err) + } + bdm := f.lastLTData.BlockDeviceMappings + if len(bdm) != 1 || awssdk.ToString(bdm[0].DeviceName) != "/dev/xvda" || bdm[0].Ebs == nil { + t.Fatalf("BlockDeviceMappings = %+v, want one root mapping on /dev/xvda", bdm) + } + ebs := bdm[0].Ebs + if awssdk.ToInt32(ebs.VolumeSize) != tc.wantSize || ebs.VolumeType != ec2types.VolumeTypeGp3 || + !awssdk.ToBool(ebs.DeleteOnTermination) { + t.Fatalf("root EBS = size %d type %q delete %v; want %d gp3 true", + awssdk.ToInt32(ebs.VolumeSize), ebs.VolumeType, awssdk.ToBool(ebs.DeleteOnTermination), tc.wantSize) + } + }) + } +} diff --git a/pkg/provider/catalog/data/pricing.go b/pkg/provider/catalog/data/pricing.go index ddf4459..3265563 100644 --- a/pkg/provider/catalog/data/pricing.go +++ b/pkg/provider/catalog/data/pricing.go @@ -62,3 +62,18 @@ func ModalCPUCostPerHour(cpuCores float64) float64 { func ModalMemoryCostPerHour(memoryMiB int) float64 { return float64(memoryMiB) / mibPerGiB * ModalMemoryPricePerGiBHour } + +// AWSGP3PricePerGBHour is gp3's US East (N. Virginia) $0.08/GB-month from +// aws.amazon.com/ebs/pricing (2026-10-04), spread over an average month. One rate for +// every region, since the catalog has no region axis (see provider.PriceRequest); other +// regions differ by a few cents per GB-month. +const AWSGP3PricePerGBHour = 0.08 / hoursPerMonth + +// hoursPerMonth is 365 days / 12, the conversion for a rate quoted per month. +const hoursPerMonth = 730 + +// AWSRootVolumeCostPerHour is what EBS charges for a gp3 root volume of diskGiB, to be ADDED +// to the instance price, which covers no storage. +func AWSRootVolumeCostPerHour(diskGiB int) float64 { + return float64(diskGiB) * AWSGP3PricePerGBHour +} diff --git a/pkg/provider/modal/modal.go b/pkg/provider/modal/modal.go index dbc72a8..995700a 100644 --- a/pkg/provider/modal/modal.go +++ b/pkg/provider/modal/modal.go @@ -395,6 +395,10 @@ func (p *Provider) ResolveRegions(declared, narrowTo []string) []string { return []string{strings.Join(regions, regionSeparator)} } +// modalFreeDiskGiB is the per-container disk Modal grants by default, without charge, and the +// most Nebula can get: sandboxSpecFromPod refuses a Pod asking for more. +const modalFreeDiskGiB = 512 + // PricePerHour overrides catalog.Base's all-in reading of the catalog, because Modal // meters CPU and memory SEPARATELY from the accelerator: a modal.csv row prices ONE GPU // and nothing else, so the sandbox's real rate is that plus what its reservation costs. @@ -408,7 +412,14 @@ func (p *Provider) ResolveRegions(declared, narrowTo []string) []string { // applies its own defaults, and we do not know them. Unpriced is the honest answer — a 0 // would be read as free. A GPU sandbox in that state still prices, understating by those // same defaults, which is immaterial beside the accelerator. +// +// Disk adds nothing up to modalFreeDiskGiB; above it is ErrNoPrice, since such a Pod is +// never launched (see sandboxSpecFromPod). func (p *Provider) PricePerHour(req provider.PriceRequest) (float64, error) { + if req.DiskGiB > modalFreeDiskGiB { + return 0, fmt.Errorf("modal: %d GiB disk exceeds the unbilled %d GiB: %w", + req.DiskGiB, modalFreeDiskGiB, provider.ErrNoPrice) + } metered := data.ModalCPUCostPerHour(req.CPUCores) + data.ModalMemoryCostPerHour(req.MemoryMiB) if req.AcceleratorType == "" { @@ -610,6 +621,12 @@ func (p *Provider) sandboxSpecFromPod(pod *corev1.Pod, req provider.ProvisionReq } c := pod.Spec.Containers[0] + // Refused, not launched short: Nebula cannot ask Modal for more than its default disk + // (the SDK has no field for it), so the workload's writes past modalFreeDiskGiB would fail. + if diskGiB := util.PodEphemeralStorageGiB(pod); diskGiB > modalFreeDiskGiB { + return SandboxSpec{}, fmt.Errorf("modal: %d GiB disk exceeds the %d GiB Modal provides", diskGiB, modalFreeDiskGiB) + } + tags := map[string]string{ClaimTagKey: req.ClaimName} // Record probe-ness alongside identity so observe can recover it later; see // ProbeTagKey for why this cannot be re-derived at observation time. The tag diff --git a/pkg/provider/modal/modal_test.go b/pkg/provider/modal/modal_test.go index 2be983f..eacc966 100644 --- a/pkg/provider/modal/modal_test.go +++ b/pkg/provider/modal/modal_test.go @@ -519,6 +519,30 @@ func TestProvision_UnsupportedAccelerator(t *testing.T) { } } +func TestProvision_DiskAboveDefaultRefused(t *testing.T) { + for name, tc := range map[string]struct { + limit string + wantErr bool + }{ + "at the default": {"512Gi", false}, + "above the default": {"513Gi", true}, + } { + t.Run(name, func(t *testing.T) { + f := &fakeClient{} + p := newTestProvider(f) + pod := gpuPod("claim-disk", "H100", 1) + pod.Spec.Containers[0].Resources.Limits[corev1.ResourceEphemeralStorage] = resource.MustParse(tc.limit) + _, err := p.Provision(context.Background(), pod, provider.ProvisionRequest{ClaimName: "claim-disk"}) + if (err != nil) != tc.wantErr { + t.Fatalf("Provision err = %v, wantErr %v", err, tc.wantErr) + } + if tc.wantErr && f.createCnt != 0 { + t.Fatalf("CreateSandbox called %d times for a refused disk", f.createCnt) + } + }) + } +} + func TestClassifyProvisionError(t *testing.T) { p := newTestProvider(&fakeClient{}) denyAll := provider.BlockScope{DenyAll: true} @@ -1974,6 +1998,10 @@ func TestPricePerHour_AddsCPUAndMemory(t *testing.T) { provider.PriceRequest{CPUCores: 4, MemoryMiB: 8192}, cpuAndMem, }, + "disk within the unbilled quota adds nothing": { + provider.PriceRequest{CPUCores: 4, MemoryMiB: 8192, DiskGiB: modalFreeDiskGiB}, + cpuAndMem, + }, } for name, tc := range cases { t.Run(name, func(t *testing.T) { @@ -2000,6 +2028,10 @@ func TestPricePerHour_NoPrice(t *testing.T) { AcceleratorType: "TPU-v4", Count: 1, CapacityType: nebulav1alpha1.CapacityOnDemand, CPUCores: 4, MemoryMiB: 8192, }, + "disk above the unbilled quota": { + AcceleratorType: "H100", Count: 1, CapacityType: nebulav1alpha1.CapacityOnDemand, + CPUCores: 4, MemoryMiB: 8192, DiskGiB: modalFreeDiskGiB + 1, + }, } { t.Run(name, func(t *testing.T) { got, err := p.PricePerHour(req) diff --git a/pkg/provider/pricing.go b/pkg/provider/pricing.go index 1ccdec0..5db1a7f 100644 --- a/pkg/provider/pricing.go +++ b/pkg/provider/pricing.go @@ -56,6 +56,10 @@ type PriceRequest struct { // instance price is all-in. CPUCores float64 MemoryMiB int + // DiskGiB is the workload's disk (see util.PodEphemeralStorageGiB), 0 when unset. Priced + // only where disk is a billed line of its own; Modal ignores it, since its default disk + // is unbilled and Nebula never requests more. + DiskGiB int } // Pricer reports the hourly USD cost of what a Provision would create: ONE rate, with every diff --git a/pkg/util/resources.go b/pkg/util/resources.go index 2aba96e..c798ba9 100644 --- a/pkg/util/resources.go +++ b/pkg/util/resources.go @@ -59,3 +59,32 @@ func reservedQty(c *corev1.Container, name corev1.ResourceName) resource.Quantit } return resource.Quantity{} } + +// gibBytes is one GiB, the unit provider.PriceRequest quotes disk in. +const gibBytes = 1024 * mibBytes + +// PodEphemeralStorageGiB reads the first container's ephemeral-storage as whole GiB, rounded +// UP, or 0 when unset. Limit first, unlike PodReservation: a provisioned disk is a hard cap, +// so sizing it to the request would fail writes the limit entitles the workload to. Unbounded: +// each provider enforces its own maximum. +func PodEphemeralStorageGiB(pod *corev1.Pod) int { + if pod == nil || len(pod.Spec.Containers) == 0 { + return 0 + } + c := &pod.Spec.Containers[0] + q, ok := c.Resources.Limits[corev1.ResourceEphemeralStorage] + if !ok { + q = c.Resources.Requests[corev1.ResourceEphemeralStorage] + } + if q.Sign() <= 0 { + return 0 + } + // Not (v + gibBytes - 1) / gibBytes: Value saturates at MaxInt64 for an admitted but huge + // Quantity, and the addition would wrap that negative. + v := q.Value() + gib := v / gibBytes + if v%gibBytes != 0 { + gib++ + } + return int(gib) +} diff --git a/pkg/util/resources_test.go b/pkg/util/resources_test.go index e7909bb..ea5959e 100644 --- a/pkg/util/resources_test.go +++ b/pkg/util/resources_test.go @@ -17,6 +17,7 @@ limitations under the License. package util import ( + "math" "testing" corev1 "k8s.io/api/core/v1" @@ -107,3 +108,30 @@ func TestPodReservation_FirstContainerOnly(t *testing.T) { t.Fatalf("PodReservation cpu = %v, want 1 (the first container's)", cpu) } } + +func TestPodEphemeralStorageGiB(t *testing.T) { + storage := func(s string) corev1.ResourceList { + return corev1.ResourceList{corev1.ResourceEphemeralStorage: resource.MustParse(s)} + } + cases := map[string]struct { + pod *corev1.Pod + want int + }{ + "limit wins over request": {podWith(storage("10Gi"), storage("50Gi")), 50}, + "request alone": {podWith(storage("10Gi"), nil), 10}, + "rounds up to whole GiB": {podWith(storage("1500Mi"), nil), 2}, + "decimal units round up": {podWith(storage("20G"), nil), 19}, + // Value saturates at MaxInt64 bytes; rounding must not wrap it negative. + "near MaxInt64 bytes": {podWith(storage("1Gi"), storage("9223372035Gi")), math.MaxInt64/(1<<30) + 1}, + "unset": {podWith(nil, nil), 0}, + "no containers": {&corev1.Pod{}, 0}, + "nil pod": {nil, 0}, + } + for name, tc := range cases { + t.Run(name, func(t *testing.T) { + if got := PodEphemeralStorageGiB(tc.pod); got != tc.want { + t.Errorf("PodEphemeralStorageGiB = %d, want %d", got, tc.want) + } + }) + } +} diff --git a/pkg/vnode/node.go b/pkg/vnode/node.go index b537f71..0bc045d 100644 --- a/pkg/vnode/node.go +++ b/pkg/vnode/node.go @@ -399,12 +399,16 @@ const nvidiaGPUResource = "nvidia.com/gpu" // (amd.com/gpu, a typed MIG key) would have its GPU Pods rejected before provisioning; // when one lands, have Capabilities declare its resource keys and build capacity from // that. -// Support at least 1k workloads with 64 cpu, 256 Gib, 8GPU per virtual node. +// Support at least 1k workloads with 64 cpu, 256 Gib, 8GPU, 16 TiB disk per virtual node. +// ephemeral-storage must be listed: an unadvertised resource is zero allocatable, so every +// Pod sizing its disk (util.PodEphemeralStorageGiB) would fail to schedule. 16 TiB is gp3's +// largest volume. func virtualCapacity() corev1.ResourceList { return corev1.ResourceList{ - corev1.ResourceCPU: resource.MustParse("64k"), - corev1.ResourceMemory: resource.MustParse("250Ti"), - corev1.ResourcePods: resource.MustParse("1k"), - nvidiaGPUResource: resource.MustParse("8k"), + corev1.ResourceCPU: resource.MustParse("64k"), + corev1.ResourceMemory: resource.MustParse("250Ti"), + corev1.ResourceEphemeralStorage: resource.MustParse("16Pi"), + corev1.ResourcePods: resource.MustParse("1k"), + nvidiaGPUResource: resource.MustParse("8k"), } } diff --git a/pkg/vnode/node_test.go b/pkg/vnode/node_test.go index 17fa3e3..a82cfa3 100644 --- a/pkg/vnode/node_test.go +++ b/pkg/vnode/node_test.go @@ -27,6 +27,7 @@ import ( "github.com/virtual-kubelet/virtual-kubelet/errdefs" vknode "github.com/virtual-kubelet/virtual-kubelet/node" corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" "k8s.io/client-go/informers" @@ -266,3 +267,32 @@ func firstEnvtestBinaryDir() string { } return "" } + +// TestVirtualCapacity_FitsTheScaleTarget replays the scheduler's resource-fit check: every +// resource a workload may request must be allocatable, at 1k workloads of the per-workload +// maximum, or those Pods are rejected before any provider sees them. +func TestVirtualCapacity_FitsTheScaleTarget(t *testing.T) { + const workloads = 1000 + perWorkload := corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse("64"), + corev1.ResourceMemory: resource.MustParse("256Gi"), + corev1.ResourceEphemeralStorage: resource.MustParse("16Ti"), + nvidiaGPUResource: resource.MustParse("8"), + } + allocatable := virtualCapacity() + for name, q := range perWorkload { + avail, ok := allocatable[name] + if !ok { + t.Errorf("%s is not advertised, so any Pod requesting it cannot schedule", name) + continue + } + total := q.DeepCopy() + total.Mul(workloads) + if avail.Cmp(total) < 0 { + t.Errorf("%s allocatable %s < %d workloads x %s", name, avail.String(), workloads, q.String()) + } + } + if pods := allocatable[corev1.ResourcePods]; pods.Value() < workloads { + t.Errorf("pods allocatable %d < %d", pods.Value(), workloads) + } +}