wilfred-s commented on code in PR #1053:
URL: https://github.com/apache/yunikorn-k8shim/pull/1053#discussion_r3620126212


##########
pkg/cache/utils.go:
##########
@@ -58,5 +61,77 @@ func GetTaskGroupsFromAnnotation(pod *v1.Pod) ([]TaskGroup, 
error) {
                                taskGroupInfo)
                }
        }
+       if err := validateTaskGroupResources(taskGroups); err != nil {
+               return nil, err
+       }

Review Comment:
   Instead of doing a new range can this not be included in the range above? 
Saves one loop
   ```
   totals := make(map[string]int64)
   for range taskGroups {
   ...
       err := validateTaskGroupResources(taskGroup, totals)
       if err != nil {
           return nil, err
       }
   }
   ```
   Check totals inside the validate call directly.



##########
pkg/cache/utils.go:
##########
@@ -58,5 +61,77 @@ func GetTaskGroupsFromAnnotation(pod *v1.Pod) ([]TaskGroup, 
error) {
                                taskGroupInfo)
                }
        }
+       if err := validateTaskGroupResources(taskGroups); err != nil {
+               return nil, err
+       }
        return taskGroups, nil
 }
+
+// maxMilliCPU is the largest cpu quantity whose MilliValue() is representable 
as an
+// int64. It is used to reject a cpu minResource whose accessor would overflow
+// before the placeholder ask is computed.
+var maxMilliCPU = resource.NewMilliQuantity(math.MaxInt64, resource.DecimalSI)
+
+// validateTaskGroupResources rejects task groups whose placeholder ask cannot 
be
+// computed as a non-negative int64. It mirrors what GetTGResource + common.Add
+// actually emit: each task group contributes an implicit "pods" count plus one
+// entry per minResource under its scheduler-interface canonical name ("cpu" 
becomes
+// siCommon.CPU, i.e. "vcore"), and the groups are summed. minMember and 
minResource
+// come unbounded from the pod annotation, so an unchecked int64 accessor, a
+// minMember*minResource product, or the aggregate across task groups can wrap 
to a
+// negative value that is then sent to the core as PlaceholderAsk. Because 
"cpu"
+// canonicalizes to the same key as "vcore" and an explicit "pods" collides 
with the
+// implicit count, the aggregate is computed in the canonical key-space with
+// same-key collisions rejected, so the check matches the real placeholder ask.
+func validateTaskGroupResources(taskGroups []TaskGroup) error {
+       totals := make(map[string]int64)
+       for _, taskGroup := range taskGroups {
+               members := int64(taskGroup.MinMember)
+               // group mirrors the si.Resource GetTGResource emits for this 
task group.
+               group := map[string]int64{"pods": members}

Review Comment:
   update totals directly with a check for overflow and return error on fail.



##########
pkg/cache/utils.go:
##########
@@ -58,5 +61,77 @@ func GetTaskGroupsFromAnnotation(pod *v1.Pod) ([]TaskGroup, 
error) {
                                taskGroupInfo)
                }
        }
+       if err := validateTaskGroupResources(taskGroups); err != nil {
+               return nil, err
+       }
        return taskGroups, nil
 }
+
+// maxMilliCPU is the largest cpu quantity whose MilliValue() is representable 
as an
+// int64. It is used to reject a cpu minResource whose accessor would overflow
+// before the placeholder ask is computed.
+var maxMilliCPU = resource.NewMilliQuantity(math.MaxInt64, resource.DecimalSI)
+
+// validateTaskGroupResources rejects task groups whose placeholder ask cannot 
be
+// computed as a non-negative int64. It mirrors what GetTGResource + common.Add
+// actually emit: each task group contributes an implicit "pods" count plus one
+// entry per minResource under its scheduler-interface canonical name ("cpu" 
becomes
+// siCommon.CPU, i.e. "vcore"), and the groups are summed. minMember and 
minResource
+// come unbounded from the pod annotation, so an unchecked int64 accessor, a
+// minMember*minResource product, or the aggregate across task groups can wrap 
to a
+// negative value that is then sent to the core as PlaceholderAsk. Because 
"cpu"
+// canonicalizes to the same key as "vcore" and an explicit "pods" collides 
with the
+// implicit count, the aggregate is computed in the canonical key-space with
+// same-key collisions rejected, so the check matches the real placeholder ask.
+func validateTaskGroupResources(taskGroups []TaskGroup) error {

Review Comment:
   Change signature, see above



##########
pkg/cache/utils.go:
##########
@@ -58,5 +61,77 @@ func GetTaskGroupsFromAnnotation(pod *v1.Pod) ([]TaskGroup, 
error) {
                                taskGroupInfo)
                }
        }
+       if err := validateTaskGroupResources(taskGroups); err != nil {
+               return nil, err
+       }
        return taskGroups, nil
 }
+
+// maxMilliCPU is the largest cpu quantity whose MilliValue() is representable 
as an
+// int64. It is used to reject a cpu minResource whose accessor would overflow
+// before the placeholder ask is computed.
+var maxMilliCPU = resource.NewMilliQuantity(math.MaxInt64, resource.DecimalSI)
+
+// validateTaskGroupResources rejects task groups whose placeholder ask cannot 
be
+// computed as a non-negative int64. It mirrors what GetTGResource + common.Add
+// actually emit: each task group contributes an implicit "pods" count plus one
+// entry per minResource under its scheduler-interface canonical name ("cpu" 
becomes
+// siCommon.CPU, i.e. "vcore"), and the groups are summed. minMember and 
minResource
+// come unbounded from the pod annotation, so an unchecked int64 accessor, a
+// minMember*minResource product, or the aggregate across task groups can wrap 
to a
+// negative value that is then sent to the core as PlaceholderAsk. Because 
"cpu"
+// canonicalizes to the same key as "vcore" and an explicit "pods" collides 
with the
+// implicit count, the aggregate is computed in the canonical key-space with
+// same-key collisions rejected, so the check matches the real placeholder ask.
+func validateTaskGroupResources(taskGroups []TaskGroup) error {
+       totals := make(map[string]int64)
+       for _, taskGroup := range taskGroups {
+               members := int64(taskGroup.MinMember)
+               // group mirrors the si.Resource GetTGResource emits for this 
task group.
+               group := map[string]int64{"pods": members}
+               for resName, quantity := range taskGroup.MinResource {
+                       value, err := checkedResourceValue(resName, quantity, 
taskGroup.Name)
+                       if err != nil {
+                               return err
+                       }
+                       canonical := resName
+                       if resName == v1.ResourceCPU.String() {
+                               canonical = siCommon.CPU
+                       }
+                       if value != 0 && members > math.MaxInt64/value {
+                               return fmt.Errorf("minMember %d times 
minResource %q overflows int64 in taskGroup %q", members, resName, 
taskGroup.Name)
+                       }
+                       if _, exists := group[canonical]; exists {
+                               return fmt.Errorf("minResource %q in taskGroup 
%q collides with canonical resource %q", resName, taskGroup.Name, canonical)
+                       }
+                       group[canonical] = members * value
+               }
+               for name, value := range group {
+                       if totals[name] > math.MaxInt64-value {
+                               return fmt.Errorf("aggregate placeholder 
request for %q overflows int64 across taskGroups", name)
+                       }
+                       totals[name] += value
+               }
+       }
+       return nil
+}
+
+// checkedResourceValue returns the int64 value GetTGResource uses for a
+// (resourceName, quantity) pair, rejecting a negative or 
non-int64-representable
+// quantity. cpu is measured in milli-units (MilliValue), everything else in 
whole
+// units (Value), matching GetTGResource.
+func checkedResourceValue(resName string, quantity resource.Quantity, 
groupName string) (int64, error) {

Review Comment:
   This can be done directly in the loop above



##########
pkg/cache/utils.go:
##########
@@ -58,5 +61,77 @@ func GetTaskGroupsFromAnnotation(pod *v1.Pod) ([]TaskGroup, 
error) {
                                taskGroupInfo)
                }
        }
+       if err := validateTaskGroupResources(taskGroups); err != nil {
+               return nil, err
+       }
        return taskGroups, nil
 }
+
+// maxMilliCPU is the largest cpu quantity whose MilliValue() is representable 
as an
+// int64. It is used to reject a cpu minResource whose accessor would overflow
+// before the placeholder ask is computed.
+var maxMilliCPU = resource.NewMilliQuantity(math.MaxInt64, resource.DecimalSI)
+
+// validateTaskGroupResources rejects task groups whose placeholder ask cannot 
be
+// computed as a non-negative int64. It mirrors what GetTGResource + common.Add
+// actually emit: each task group contributes an implicit "pods" count plus one
+// entry per minResource under its scheduler-interface canonical name ("cpu" 
becomes
+// siCommon.CPU, i.e. "vcore"), and the groups are summed. minMember and 
minResource
+// come unbounded from the pod annotation, so an unchecked int64 accessor, a
+// minMember*minResource product, or the aggregate across task groups can wrap 
to a
+// negative value that is then sent to the core as PlaceholderAsk. Because 
"cpu"
+// canonicalizes to the same key as "vcore" and an explicit "pods" collides 
with the
+// implicit count, the aggregate is computed in the canonical key-space with
+// same-key collisions rejected, so the check matches the real placeholder ask.
+func validateTaskGroupResources(taskGroups []TaskGroup) error {
+       totals := make(map[string]int64)
+       for _, taskGroup := range taskGroups {
+               members := int64(taskGroup.MinMember)
+               // group mirrors the si.Resource GetTGResource emits for this 
task group.
+               group := map[string]int64{"pods": members}
+               for resName, quantity := range taskGroup.MinResource {
+                       value, err := checkedResourceValue(resName, quantity, 
taskGroup.Name)
+                       if err != nil {
+                               return err
+                       }
+                       canonical := resName
+                       if resName == v1.ResourceCPU.String() {
+                               canonical = siCommon.CPU
+                       }
+                       if value != 0 && members > math.MaxInt64/value {
+                               return fmt.Errorf("minMember %d times 
minResource %q overflows int64 in taskGroup %q", members, resName, 
taskGroup.Name)
+                       }
+                       if _, exists := group[canonical]; exists {
+                               return fmt.Errorf("minResource %q in taskGroup 
%q collides with canonical resource %q", resName, taskGroup.Name, canonical)
+                       }
+                       group[canonical] = members * value
+               }
+               for name, value := range group {
+                       if totals[name] > math.MaxInt64-value {
+                               return fmt.Errorf("aggregate placeholder 
request for %q overflows int64 across taskGroups", name)
+                       }
+                       totals[name] += value
+               }

Review Comment:
   directly in the above loop



##########
pkg/cache/utils.go:
##########
@@ -58,5 +61,77 @@ func GetTaskGroupsFromAnnotation(pod *v1.Pod) ([]TaskGroup, 
error) {
                                taskGroupInfo)
                }
        }
+       if err := validateTaskGroupResources(taskGroups); err != nil {
+               return nil, err
+       }
        return taskGroups, nil
 }
+
+// maxMilliCPU is the largest cpu quantity whose MilliValue() is representable 
as an
+// int64. It is used to reject a cpu minResource whose accessor would overflow
+// before the placeholder ask is computed.
+var maxMilliCPU = resource.NewMilliQuantity(math.MaxInt64, resource.DecimalSI)
+
+// validateTaskGroupResources rejects task groups whose placeholder ask cannot 
be
+// computed as a non-negative int64. It mirrors what GetTGResource + common.Add
+// actually emit: each task group contributes an implicit "pods" count plus one
+// entry per minResource under its scheduler-interface canonical name ("cpu" 
becomes
+// siCommon.CPU, i.e. "vcore"), and the groups are summed. minMember and 
minResource
+// come unbounded from the pod annotation, so an unchecked int64 accessor, a
+// minMember*minResource product, or the aggregate across task groups can wrap 
to a
+// negative value that is then sent to the core as PlaceholderAsk. Because 
"cpu"
+// canonicalizes to the same key as "vcore" and an explicit "pods" collides 
with the
+// implicit count, the aggregate is computed in the canonical key-space with
+// same-key collisions rejected, so the check matches the real placeholder ask.
+func validateTaskGroupResources(taskGroups []TaskGroup) error {
+       totals := make(map[string]int64)
+       for _, taskGroup := range taskGroups {
+               members := int64(taskGroup.MinMember)
+               // group mirrors the si.Resource GetTGResource emits for this 
task group.
+               group := map[string]int64{"pods": members}
+               for resName, quantity := range taskGroup.MinResource {
+                       value, err := checkedResourceValue(resName, quantity, 
taskGroup.Name)

Review Comment:
   * check sign here first, error on fail
   * make canonical for CPU 
   * calculate max value based on type (cpu/vcore or not) and members:
   ```
   milliConvert := 1
   if cpu {
        milliConvert = 1000
   }
   maxVal := math.MaxInt64/(members*milliConvert)
   if quantity.CmpInt64(maxVal) > 0 {
    return err...
   }
   ```
   * update totals directly with a check for overflow and return error on fail



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to