Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
7 changes: 7 additions & 0 deletions internal/common/resource/resource.go
Original file line number Diff line number Diff line change
Expand Up @@ -264,6 +264,13 @@ func TotalPodResourceRequest(podSpec *v1.PodSpec) ComputeResources {
totalResources.Max(containerResource)
}
}

// Pod-level resources (KEP-2837): the effective request is
// max(sum of container requests, pod-level request) per resource. Inert when
// podSpec.Resources is unset.
if podSpec.Resources != nil {
totalResources.Max(FromResourceList(podSpec.Resources.Requests))
}
return totalResources
}

Expand Down
43 changes: 43 additions & 0 deletions internal/common/resource/resource_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,49 @@ func TestTotalResourceRequest_ShouldCombineMaxInitContainerResourcesWithSummedCo
assert.Equal(t, result, FromResourceList(expectedResult))
}

// Pod-level resources (KEP-2837): TotalPodResourceRequest uses the effective
// request max(sum(containers), pod-level).
func TestTotalResourceRequest_PodLevelResources(t *testing.T) {
Comment thread
washcycle marked this conversation as resolved.
t.Run("pod-level only (empty containers) is accounted at the pod-level value", func(t *testing.T) {
podLevel := makeContainerResource(4, 16)
pod := makePodWithResource([]*v1.ResourceList{}, []*v1.ResourceList{})
pod.Spec.Resources = &v1.ResourceRequirements{Requests: podLevel, Limits: podLevel}

result := TotalPodResourceRequest(&pod.Spec)
assert.Equal(t, FromResourceList(makeContainerResource(4, 16)), result)
})

t.Run("effective request is max of container-sum and pod-level", func(t *testing.T) {
container := makeContainerResource(2, 4) // sum of one container = 2cpu/4Gi
podLevel := makeContainerResource(4, 2) // pod-level = 4cpu/2Gi
pod := makePodWithResource([]*v1.ResourceList{&container}, []*v1.ResourceList{})
pod.Spec.Resources = &v1.ResourceRequirements{Requests: podLevel, Limits: podLevel}

// cpu: max(2, 4) = 4 ; memory: max(4, 2) = 4
result := TotalPodResourceRequest(&pod.Spec)
assert.Equal(t, FromResourceList(makeContainerResource(4, 4)), result)
})

t.Run("container-sum wins when it exceeds pod-level", func(t *testing.T) {
container := makeContainerResource(8, 8)
podLevel := makeContainerResource(4, 4)
pod := makePodWithResource([]*v1.ResourceList{&container}, []*v1.ResourceList{})
pod.Spec.Resources = &v1.ResourceRequirements{Requests: podLevel, Limits: podLevel}

result := TotalPodResourceRequest(&pod.Spec)
assert.Equal(t, FromResourceList(makeContainerResource(8, 8)), result)
})

t.Run("nil pod-level leaves upstream behaviour unchanged", func(t *testing.T) {
container := makeContainerResource(2, 4)
pod := makePodWithResource([]*v1.ResourceList{&container}, []*v1.ResourceList{})
// pod.Spec.Resources stays nil

result := TotalPodResourceRequest(&pod.Spec)
assert.Equal(t, FromResourceList(makeContainerResource(2, 4)), result)
})
}

func TestTotalResourceRequest_NativeSidecarsShouldBeSummed(t *testing.T) {
mainResource := makeContainerResource(2, 1)
sidecarResource := makeContainerResource(1, 1)
Expand Down
19 changes: 11 additions & 8 deletions internal/lookoutingester/instructions/instructions.go
Original file line number Diff line number Diff line change
Expand Up @@ -535,18 +535,21 @@ func getJobResources(job *api.Job) jobResources {

podSpec := job.GetMainPodSpec()

for _, container := range podSpec.Containers {
resources.Cpu += getResource(container, v1.ResourceCPU, true)
resources.Memory += getResource(container, v1.ResourceMemory, false)
resources.EphemeralStorage += getResource(container, v1.ResourceEphemeralStorage, false)
resources.Gpu += getResource(container, "nvidia.com/gpu", false)
}
// Use the canonical effective-request computation so the reported footprint
// matches what the scheduler bills: sum of main containers + native sidecars,
// max over classic init containers, and max with the pod-level block (KEP-2837).
// This also fixes the prior undercount that summed only main containers.
requests := api.SchedulingResourceRequirementsFromPodSpec(podSpec).Requests
resources.Cpu = getResourceFromList(requests, v1.ResourceCPU, true)
resources.Memory = getResourceFromList(requests, v1.ResourceMemory, false)
resources.EphemeralStorage = getResourceFromList(requests, v1.ResourceEphemeralStorage, false)
resources.Gpu = getResourceFromList(requests, "nvidia.com/gpu", false)

return resources
}

func getResource(container v1.Container, resourceName v1.ResourceName, useMillis bool) int64 {
resource, ok := container.Resources.Requests[resourceName]
func getResourceFromList(rl v1.ResourceList, resourceName v1.ResourceName, useMillis bool) int64 {
resource, ok := rl[resourceName]
if !ok {
return 0
}
Expand Down
6 changes: 6 additions & 0 deletions internal/scheduler/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -231,6 +231,12 @@ func (srv *ExecutorApi) dropDisallowedResources(pod *v1.PodSpec) {
}
srv.dropDisallowedResourcesFromContainers(pod.InitContainers)
srv.dropDisallowedResourcesFromContainers(pod.Containers)
// Pod-level resources (KEP-2837) must be filtered by the same allow-list, else
// a disallowed resource could bypass it via the pod-level block.
if pod.Resources != nil {
removeDisallowedKeys(pod.Resources.Limits, srv.allowedResources)
removeDisallowedKeys(pod.Resources.Requests, srv.allowedResources)
}
}
Comment thread
greptile-apps[bot] marked this conversation as resolved.

func (srv *ExecutorApi) dropDisallowedResourcesFromContainers(containers []v1.Container) {
Expand Down
5 changes: 5 additions & 0 deletions internal/server/configuration/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,11 @@ type SubmissionConfig struct {
AddGangIdLabel bool
// Controls whether custom service names are allowed
AllowCustomServiceNames bool
// When true, honour Kubernetes pod-level resources (KEP-2837, podSpec.Resources):
// a container may omit its own resources if the pod-level block is set, and
// accounting uses max(sum(container requests), podLevel.Requests). Default false
// preserves container-only behaviour.
PodLevelResources bool
}

// TODO: we can probably just typedef this to map[string]string
Expand Down
64 changes: 64 additions & 0 deletions internal/server/submit/validation/submit_request.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"strings"

"github.com/pkg/errors"
v1 "k8s.io/api/core/v1"
"k8s.io/component-helpers/scheduling/corev1/nodeaffinity"

"github.com/armadaproject/armada/internal/common/constants"
Expand Down Expand Up @@ -252,8 +253,21 @@ func validateResources(j *api.JobSubmitRequestItem, config configuration.Submiss
if maxOversubscriptionByResource == nil {
maxOversubscriptionByResource = map[string]float64{}
}
// Pod-level resources (KEP-2837): when enabled and a pod-level block is set, a
// container may omit its own resources. Containers that do set resources are
// still validated below, as is the pod-level block itself.
podLevelResourcesEnabled := config.PodLevelResources && spec.Resources != nil
if podLevelResourcesEnabled {
if err := validatePodLevelResources(spec.Resources, maxOversubscriptionByResource, config); err != nil {
return err
}
}
for _, container := range armadaslices.Concatenate(spec.Containers, spec.InitContainers) {
if len(container.Resources.Requests) == 0 && len(container.Resources.Limits) == 0 {
if podLevelResourcesEnabled {
// Budget is supplied at the pod level; nothing to validate for this container.
continue
}
return fmt.Errorf("container %v has no resources specified", container.Name)
}

Expand Down Expand Up @@ -308,6 +322,56 @@ func validateResources(j *api.JobSubmitRequestItem, config configuration.Submiss
return nil
}

// validatePodLevelResources validates a pod-level resources block (KEP-2837),
// mirroring the per-container checks: requests and limits must be non-negative,
// cover the same resource set, satisfy limit >= request within the
// max-oversubscription ratio, and meet MinJobResources.
func validatePodLevelResources(
resources *v1.ResourceRequirements,
maxOversubscriptionByResource map[string]float64,
config configuration.SubmissionConfig,
) error {
if len(resources.Requests) == 0 && len(resources.Limits) == 0 {
return fmt.Errorf("pod-level resources block is empty")
}
if len(resources.Requests) != len(resources.Limits) {
return fmt.Errorf("pod-level resources define different resources for requests and limits")
}
for resourceName, request := range resources.Requests {
if request.Sign() < 0 {
return fmt.Errorf("pod-level resources define negative request (%s) for resource %s", request.String(), resourceName)
}
}
for resourceName, limit := range resources.Limits {
if limit.Sign() < 0 {
return fmt.Errorf("pod-level resources define negative limit (%s) for resource %s", limit.String(), resourceName)
}
}
for resourceName, request := range resources.Requests {
limit, ok := resources.Limits[resourceName]
if !ok {
return fmt.Errorf("pod-level resources define %s for requests but not limits", resourceName)
}
if limit.MilliValue() < request.MilliValue() {
return fmt.Errorf("pod-level resources define %s with limits smaller than requests", resourceName)
}
maxOversubscription, ok := maxOversubscriptionByResource[resourceName.String()]
if !ok {
maxOversubscription = 1.0
}
if float64(limit.MilliValue()) > maxOversubscription*float64(request.MilliValue()) {
return fmt.Errorf("pod-level resources define %s with limits greater than %.2f*requests", resourceName, maxOversubscription)
}
}
for rc, podRsc := range resources.Requests {
serverRsc, nonEmpty := config.MinJobResources[rc]
if nonEmpty && podRsc.Value() < serverRsc.Value() {
return fmt.Errorf("pod-level %s requests (%s) below server minimum (%s)", rc, &podRsc, &serverRsc)
}
}
Comment thread
greptile-apps[bot] marked this conversation as resolved.
Outdated
return nil
}

// jobAdapter turns JobSubmitRequestItem into a MinimalJob
// This is needed for gang information to be extracted
type jobAdapter struct {
Expand Down
77 changes: 77 additions & 0 deletions internal/server/submit/validation/submit_request_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1150,6 +1150,83 @@ func TestValidateResources(t *testing.T) {
}
}

// TestValidateResources_PodLevel covers Kubernetes pod-level resources (KEP-2837),
// gated by SubmissionConfig.PodLevelResources.
func TestValidateResources_PodLevel(t *testing.T) {
oneCpu := v1.ResourceList{v1.ResourceCPU: resource.MustParse("1")}
twoCpu := v1.ResourceList{v1.ResourceCPU: resource.MustParse("2")}
negativeCpu := v1.ResourceList{v1.ResourceCPU: resource.MustParse("-1")}

// Container that declares no resources of its own.
emptyContainer := []v1.Container{{Name: "main"}}

req := func(podResources *v1.ResourceRequirements, containers []v1.Container) *api.JobSubmitRequestItem {
return &api.JobSubmitRequestItem{
PodSpec: &v1.PodSpec{
Containers: containers,
Resources: podResources,
},
}
}

tests := map[string]struct {
req *api.JobSubmitRequestItem
podLevelEnabled bool
expectSuccess bool
}{
"pod-level-only accepted when feature enabled": {
req: req(&v1.ResourceRequirements{Requests: oneCpu, Limits: oneCpu}, emptyContainer),
podLevelEnabled: true,
expectSuccess: true,
},
"pod-level-only rejected when feature disabled": {
req: req(&v1.ResourceRequirements{Requests: oneCpu, Limits: oneCpu}, emptyContainer),
podLevelEnabled: false,
expectSuccess: false,
},
"empty container with NO pod-level still rejected even when enabled": {
req: req(nil, emptyContainer),
podLevelEnabled: true,
expectSuccess: false,
},
"pod-level with limits < requests rejected": {
req: req(&v1.ResourceRequirements{Requests: twoCpu, Limits: oneCpu}, emptyContainer),
podLevelEnabled: true,
expectSuccess: false,
},
"pod-level negative request rejected": {
req: req(&v1.ResourceRequirements{Requests: negativeCpu, Limits: negativeCpu}, emptyContainer),
podLevelEnabled: true,
expectSuccess: false,
},
"pod-level empty block rejected": {
req: req(&v1.ResourceRequirements{}, emptyContainer),
podLevelEnabled: true,
expectSuccess: false,
},
"container-only still valid with feature enabled": {
req: req(nil, []v1.Container{{
Name: "main",
Resources: v1.ResourceRequirements{Requests: oneCpu, Limits: oneCpu},
}}),
podLevelEnabled: true,
expectSuccess: true,
},
}

for name, tc := range tests {
t.Run(name, func(t *testing.T) {
cfg := configuration.SubmissionConfig{PodLevelResources: tc.podLevelEnabled}
err := validateResources(tc.req, cfg)
if tc.expectSuccess {
assert.NoError(t, err)
} else {
assert.Error(t, err)
}
})
}
}

func TestValidateTerminationGracePeriod(t *testing.T) {
defaultMinPeriod := 30 * time.Second
defaultMaxPeriod := 300 * time.Second
Expand Down
7 changes: 7 additions & 0 deletions pkg/api/util.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,13 @@ func SchedulingResourceRequirementsFromPodSpec(podSpec *v1.PodSpec) *v1.Resource
maxResourcesToList(rv.Limits, c.Resources.Limits)
}
}

// Pod-level resources (KEP-2837): max with the pod-level request/limit so the
// scheduler reserves the pod-level budget. Inert when podSpec.Resources is unset.
if podSpec.Resources != nil {
maxResourcesToList(rv.Requests, podSpec.Resources.Requests)
maxResourcesToList(rv.Limits, podSpec.Resources.Limits)
}
Comment thread
greptile-apps[bot] marked this conversation as resolved.
return &rv
}

Expand Down
33 changes: 33 additions & 0 deletions pkg/api/util_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -373,6 +373,39 @@ func TestSchedulingResourceRequirementsFromPodSpec(t *testing.T) {
},
},
},
// Pod-level resources (KEP-2837).
"pod-level only (empty containers) uses the pod-level value": {
input: &v1.PodSpec{
Containers: []v1.Container{{}},
Resources: &v1.ResourceRequirements{
Requests: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
Limits: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
},
},
expected: &v1.ResourceRequirements{
Requests: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
Limits: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
},
},
"pod-level is max'd with the container sum": {
input: &v1.PodSpec{
Containers: []v1.Container{{
Resources: v1.ResourceRequirements{
Requests: v1.ResourceList{"cpu": QuantityWithMilliValue(1000)},
Limits: v1.ResourceList{"cpu": QuantityWithMilliValue(1000)},
},
}},
Resources: &v1.ResourceRequirements{
Requests: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
Limits: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
},
},
// max(container-sum 1000, pod-level 4000) = 4000
expected: &v1.ResourceRequirements{
Requests: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
Limits: v1.ResourceList{"cpu": QuantityWithMilliValue(4000)},
},
},
}
for name, tc := range tests {
t.Run(name, func(t *testing.T) {
Expand Down