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


##########
pkg/scheduler/partition.go:
##########
@@ -1504,6 +1504,45 @@ func (pc *PartitionContext) removeAllocation(release 
*si.AllocationRelease) ([]*
                return nil, nil
        }
 
+       // Handle shim-initiated scheduling failure: roll the allocation back 
to a pending ask

Review Comment:
   Can we factor out this code in its own method?



##########
pkg/scheduler/objects/application.go:
##########
@@ -799,6 +810,51 @@ func (sa *Application) DeallocateAsk(allocKey string) 
(*resources.Resource, erro
        return nil, fmt.Errorf("failed to locate ask with key %s", allocKey)
 }
 
+// RollbackAllocation atomically reverts an allocated ask back to pending 
state.
+// This is used when the shim fails to bind a pod after the core has allocated 
it,
+// allowing the task to be re-scheduled on a different node.
+// The application must be in Accepted or Running state for rollback to 
succeed.
+func (sa *Application) RollbackAllocation(allocKey string) 
(*resources.Resource, error) {
+       sa.Lock()
+       defer sa.Unlock()
+
+       if !sa.IsAccepted() && !sa.IsRunning() {
+               return nil, fmt.Errorf("cannot rollback allocation %s: 
application %s is in state %s", allocKey, sa.ApplicationID, sa.CurrentState())
+       }
+
+       ask := sa.requests[allocKey]
+       if ask == nil {
+               return nil, fmt.Errorf("failed to locate ask with key %s for 
rollback", allocKey)
+       }
+       alloc := sa.allocations[allocKey]
+       if alloc == nil {
+               return nil, fmt.Errorf("failed to locate allocation with key %s 
for rollback", allocKey)
+       }

Review Comment:
   The same reference (pointer) is used in both maps as the value.
   This should always be true: `ask == alloc`
   allocations must always be a subset of requests checking allocations should 
be enough
   



##########
pkg/scheduler/partition.go:
##########
@@ -1504,6 +1504,45 @@ func (pc *PartitionContext) removeAllocation(release 
*si.AllocationRelease) ([]*
                return nil, nil
        }
 
+       // Handle shim-initiated scheduling failure: roll the allocation back 
to a pending ask
+       // so the core can re-schedule it on a different node. This bypasses 
the normal
+       // remove-and-destroy path entirely. We return nil,nil to suppress the 
echo back to
+       // the shim (the shim initiated this release; no confirmation is 
needed).
+       if release.TerminationType == 
si.TerminationType_SCHEDULING_FAILED_ON_RM {
+               queue := app.GetQueue()
+               // Retrieve node ID before rolling back (RollbackAllocation 
clears it on the ask).
+               nodeID := app.GetAllocationNodeID(allocationKey)
+               res, err := app.RollbackAllocation(allocationKey)
+               if err != nil {
+                       log.Log(log.SchedPartition).Warn("failed to rollback 
allocation",
+                               zap.String("appID", appID),
+                               zap.String("allocationKey", allocationKey),
+                               zap.Error(err))
+                       return nil, nil
+               }
+               if node := pc.GetNode(nodeID); node != nil {
+                       node.RemoveAllocation(allocationKey)
+               } else {
+                       log.Log(log.SchedPartition).Warn("node not found while 
rolling back allocation",
+                               zap.String("appID", appID),
+                               zap.String("allocationKey", allocationKey),
+                               zap.String("nodeID", nodeID))
+               }
+               if err := queue.DecAllocatedResource(res); err != nil {
+                       log.Log(log.SchedPartition).Warn("failed to release 
resources from queue during rollback",
+                               zap.String("appID", appID),
+                               zap.String("allocationKey", allocationKey),
+                               zap.Error(err))
+               }
+               pc.updateAllocationCount(-1)
+               
metrics.GetQueueMetrics(queue.GetQueuePath()).AddReleasedContainers(1)

Review Comment:
   Two points:
   * are we sure that the container was added to the allocated containers 
metric already? The metric gets increased when `UpdateAllocation()` is called 
after confirming the allocation. I think that has not happened yet.
   * if above is not true: why not use `IncReleasedContainer()` ? less overhead



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