wilfred-s commented on code in PR #1098:
URL: https://github.com/apache/yunikorn-core/pull/1098#discussion_r3502974523


##########
pkg/scheduler/partition.go:
##########
@@ -42,9 +42,12 @@ import (
        "github.com/apache/yunikorn-core/pkg/scheduler/policies"
        "github.com/apache/yunikorn-core/pkg/scheduler/ugm"
        "github.com/apache/yunikorn-core/pkg/webservice/dao"
+       siCommon "github.com/apache/yunikorn-scheduler-interface/lib/go/common"
        "github.com/apache/yunikorn-scheduler-interface/lib/go/si"
 )
 
+const allocationRollbackTag = siCommon.DomainYuniKornInternal + 
"allocation-rollback"

Review Comment:
   Move to `common.constant.go`



##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc 
*objects.Allocation) (request
                return false, true, nil
        }
 
+       if existing.IsAllocated() && !allocated && 
isAllocationRollbackRequest(alloc) {
+               log.Log(log.SchedPartition).Info("handling allocation rollback",
+                       zap.String("partitionName", pc.Name),
+                       zap.String("appID", applicationID),
+                       zap.String("allocationKey", allocationKey),
+                       zap.String("nodeID", existing.GetNodeID()))
+
+               if existingNode.RemoveAllocation(allocationKey) == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from node %s", allocationKey, existing.GetNodeID())
+               }
+               // this only removes allocation from the allocations map and 
the allocation is still present in requests map
+               removed := app.RemoveAllocation(allocationKey, 
si.TerminationType_STOPPED_BY_RM)
+               if removed == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from application %s", allocationKey, applicationID)
+               }
+
+               if err := 
queue.DecAllocatedResource(existing.GetAllocatedResource()); err != nil {
+                       log.Log(log.SchedPartition).Warn("failed to release 
resources from queue during rollback",
+                               zap.String("appID", applicationID),
+                               zap.String("allocationKey", allocationKey),
+                               zap.Error(err))
+               }
+               
metrics.GetQueueMetrics(queue.GetQueuePath()).IncReleasedContainer()
+
+               if _, err := app.DeallocateAsk(allocationKey); err != nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to deallocate 
ask %s on app %s during rollback: %w", allocationKey, applicationID, err)
+               }
+
+               if app.IsCompleting() {
+                       if err := 
app.HandleApplicationEvent(objects.RunApplication); err != nil {
+                               log.Log(log.SchedPartition).Warn("failed to 
transition application back to running after rollback",
+                                       zap.String("appID", applicationID),
+                                       zap.String("allocationKey", 
allocationKey),
+                                       zap.Error(err))
+                       }
+               }

Review Comment:
   This failure should never trigger a state change. The change is state 
agnostic: the allocation is still there now in pending form.



##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc 
*objects.Allocation) (request
                return false, true, nil
        }
 
+       if existing.IsAllocated() && !allocated && 
isAllocationRollbackRequest(alloc) {
+               log.Log(log.SchedPartition).Info("handling allocation rollback",
+                       zap.String("partitionName", pc.Name),
+                       zap.String("appID", applicationID),
+                       zap.String("allocationKey", allocationKey),
+                       zap.String("nodeID", existing.GetNodeID()))
+
+               if existingNode.RemoveAllocation(allocationKey) == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from node %s", allocationKey, existing.GetNodeID())
+               }
+               // this only removes allocation from the allocations map and 
the allocation is still present in requests map
+               removed := app.RemoveAllocation(allocationKey, 
si.TerminationType_STOPPED_BY_RM)
+               if removed == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from application %s", allocationKey, applicationID)
+               }
+
+               if err := 
queue.DecAllocatedResource(existing.GetAllocatedResource()); err != nil {
+                       log.Log(log.SchedPartition).Warn("failed to release 
resources from queue during rollback",
+                               zap.String("appID", applicationID),
+                               zap.String("allocationKey", allocationKey),
+                               zap.Error(err))
+               }
+               
metrics.GetQueueMetrics(queue.GetQueuePath()).IncReleasedContainer()
+
+               if _, err := app.DeallocateAsk(allocationKey); err != nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to deallocate 
ask %s on app %s during rollback: %w", allocationKey, applicationID, err)
+               }
+
+               if app.IsCompleting() {
+                       if err := 
app.HandleApplicationEvent(objects.RunApplication); err != nil {
+                               log.Log(log.SchedPartition).Warn("failed to 
transition application back to running after rollback",
+                                       zap.String("appID", applicationID),
+                                       zap.String("allocationKey", 
allocationKey),
+                                       zap.Error(err))
+                       }
+               }
+
+               existing.SetNodeID("")
+               existing.SetBindTime(time.Time{})
+               pc.updateAllocationCount(-1)
+               if existing.IsPlaceholder() {
+                       pc.decPhAllocationCount(1)
+               }

Review Comment:
   Placeholders cannot specify volumes. They cannot fail due to volume bind 
failures.



##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc 
*objects.Allocation) (request
                return false, true, nil
        }
 
+       if existing.IsAllocated() && !allocated && 
isAllocationRollbackRequest(alloc) {
+               log.Log(log.SchedPartition).Info("handling allocation rollback",
+                       zap.String("partitionName", pc.Name),
+                       zap.String("appID", applicationID),
+                       zap.String("allocationKey", allocationKey),
+                       zap.String("nodeID", existing.GetNodeID()))
+
+               if existingNode.RemoveAllocation(allocationKey) == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from node %s", allocationKey, existing.GetNodeID())
+               }
+               // this only removes allocation from the allocations map and 
the allocation is still present in requests map
+               removed := app.RemoveAllocation(allocationKey, 
si.TerminationType_STOPPED_BY_RM)

Review Comment:
   This is not correct: we do have two maps with all allocated and pending 
allocations but they are the same object. When you remove it here it might be 
gone from the map but it is not released for scheduling.
   
   We need a specific rollback for the app and queue:
   * decrease allocated resource on app, user and queue
   * increase pending resources on app and queue
   * remove the reference from the allocated map
   * remove the node from the allocation, unset allocated flag (must become 
false)
   * add an entry to the allocLog of the allocation for the bind failure
   
   Some of these steps are scattered around in this code but it is not 
correct/complete



##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc 
*objects.Allocation) (request
                return false, true, nil
        }
 
+       if existing.IsAllocated() && !allocated && 
isAllocationRollbackRequest(alloc) {
+               log.Log(log.SchedPartition).Info("handling allocation rollback",
+                       zap.String("partitionName", pc.Name),
+                       zap.String("appID", applicationID),
+                       zap.String("allocationKey", allocationKey),
+                       zap.String("nodeID", existing.GetNodeID()))
+
+               if existingNode.RemoveAllocation(allocationKey) == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from node %s", allocationKey, existing.GetNodeID())
+               }
+               // this only removes allocation from the allocations map and 
the allocation is still present in requests map
+               removed := app.RemoveAllocation(allocationKey, 
si.TerminationType_STOPPED_BY_RM)
+               if removed == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from application %s", allocationKey, applicationID)
+               }
+
+               if err := 
queue.DecAllocatedResource(existing.GetAllocatedResource()); err != nil {
+                       log.Log(log.SchedPartition).Warn("failed to release 
resources from queue during rollback",
+                               zap.String("appID", applicationID),
+                               zap.String("allocationKey", allocationKey),
+                               zap.Error(err))
+               }
+               
metrics.GetQueueMetrics(queue.GetQueuePath()).IncReleasedContainer()

Review Comment:
   The container is not released it never ran, this metric should not increase



##########
pkg/scheduler/partition.go:
##########
@@ -1336,9 +1339,70 @@ func (pc *PartitionContext) UpdateAllocation(alloc 
*objects.Allocation) (request
                return false, true, nil
        }
 
+       if existing.IsAllocated() && !allocated && 
isAllocationRollbackRequest(alloc) {
+               log.Log(log.SchedPartition).Info("handling allocation rollback",
+                       zap.String("partitionName", pc.Name),
+                       zap.String("appID", applicationID),
+                       zap.String("allocationKey", allocationKey),
+                       zap.String("nodeID", existing.GetNodeID()))
+
+               if existingNode.RemoveAllocation(allocationKey) == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from node %s", allocationKey, existing.GetNodeID())
+               }
+               // this only removes allocation from the allocations map and 
the allocation is still present in requests map
+               removed := app.RemoveAllocation(allocationKey, 
si.TerminationType_STOPPED_BY_RM)
+               if removed == nil {
+                       metrics.GetSchedulerMetrics().IncSchedulingError()
+                       return false, false, fmt.Errorf("failed to remove 
allocation %s from application %s", allocationKey, applicationID)
+               }

Review Comment:
   This cannot happen as `existing` cannot be nil when we get here



-- 
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