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


##########
pkg/cache/task.go:
##########
@@ -343,46 +353,68 @@ func (task *Task) postTaskPending() {
 // This routine binds the pod to the allocated node.
 // It calls K8s api to bind a pod to the assigned node, this may need some 
time,
 // so we do a delay binding, background process, to avoid blocking main 
process.
-// The result of the binding is tracked and failures are properly handled.
-// If successful, we move task to next state BOUND, otherwise we fail the task
+// Volume binding and pod binding are retried with a backoff; if they 
ultimately fail
+// the allocation is rolled back to a pending ask so the core can re-schedule 
the task
+// on a different node. On success we move the task to the next state BOUND.
 func (task *Task) postTaskAllocated() {
-       go func() {
-               // we need to obtain task's lock first,
-               // this ensures no other threads modifying task state at the 
time being
-               task.lock.Lock()
-               defer task.lock.Unlock()
+       // Snapshot the fields needed for binding before launching the 
goroutine without re-acquiring the lock.
+       // We already hold task.lock (via task.handle()) during state 
transitions.
+       pod := task.pod
+       alias := task.alias
+       nodeName := task.nodeName
+       allocationKey := task.allocationKey
 
+       go func(pod *v1.Pod, alias, nodeName, allocationKey string) {
                // post a message to indicate the pod gets its allocation
-               events.GetRecorder().Eventf(task.pod.DeepCopy(),
+               events.GetRecorder().Eventf(pod.DeepCopy(),
                        nil, v1.EventTypeNormal, "Scheduled", "Scheduled",
-                       "Successfully assigned %s to node %s", task.alias, 
task.nodeName)
+                       "Successfully assigned %s to node %s", alias, nodeName)

Review Comment:
   This is not really successful as yet, the attempt is being made but that 
might fail.
   This should be moved to the point we have done the bind and we log it in 
line L407



##########
pkg/cache/task.go:
##########
@@ -625,11 +657,20 @@ func (task *Task) failWithEvent(errorMessage, 
actionReason string) {
 // pending ask so it can be re-scheduled on a different node.
 // Must be called without holding the task lock.
 func (task *Task) rollbackOnAssumePodFailure(allocationKey, nodeID string) {
-       // Read fields needed for event posting and release request.
+       task.lock.RLock()
+       alias := task.alias
+       task.lock.RUnlock()
+       task.rollbackAllocation(allocationKey, nodeID, "AssumePodFailed",
+               fmt.Sprintf("Node assignment failed for %s on node %s, it will 
be retried", alias, nodeID))
+}
+
+// rollbackAllocation resets the task allocation state and notifies the core 
to move
+// the allocation back to a pending ask so it can be re-scheduled on a 
different node.
+// Must be called without holding the task lock.
+func (task *Task) rollbackAllocation(allocationKey, nodeID, eventReason, 
eventMsg string) {

Review Comment:
   This is now called from the locked go routine, no locking here.
   Comment is correct, code is not.



##########
pkg/cache/task.go:
##########
@@ -625,11 +657,20 @@ func (task *Task) failWithEvent(errorMessage, 
actionReason string) {
 // pending ask so it can be re-scheduled on a different node.
 // Must be called without holding the task lock.
 func (task *Task) rollbackOnAssumePodFailure(allocationKey, nodeID string) {
-       // Read fields needed for event posting and release request.
+       task.lock.RLock()
+       alias := task.alias
+       task.lock.RUnlock()

Review Comment:
   Please check locking as this is called from an unlocked location. The callee 
`rollbackAllocation` needs a write lock to be held.



##########
pkg/cache/task.go:
##########
@@ -343,46 +353,68 @@ func (task *Task) postTaskPending() {
 // This routine binds the pod to the allocated node.
 // It calls K8s api to bind a pod to the assigned node, this may need some 
time,
 // so we do a delay binding, background process, to avoid blocking main 
process.
-// The result of the binding is tracked and failures are properly handled.
-// If successful, we move task to next state BOUND, otherwise we fail the task
+// Volume binding and pod binding are retried with a backoff; if they 
ultimately fail
+// the allocation is rolled back to a pending ask so the core can re-schedule 
the task
+// on a different node. On success we move the task to the next state BOUND.
 func (task *Task) postTaskAllocated() {
-       go func() {
-               // we need to obtain task's lock first,
-               // this ensures no other threads modifying task state at the 
time being
-               task.lock.Lock()
-               defer task.lock.Unlock()
+       // Snapshot the fields needed for binding before launching the 
goroutine without re-acquiring the lock.
+       // We already hold task.lock (via task.handle()) during state 
transitions.
+       pod := task.pod
+       alias := task.alias
+       nodeName := task.nodeName
+       allocationKey := task.allocationKey
 
+       go func(pod *v1.Pod, alias, nodeName, allocationKey string) {
                // post a message to indicate the pod gets its allocation
-               events.GetRecorder().Eventf(task.pod.DeepCopy(),
+               events.GetRecorder().Eventf(pod.DeepCopy(),
                        nil, v1.EventTypeNormal, "Scheduled", "Scheduled",
-                       "Successfully assigned %s to node %s", task.alias, 
task.nodeName)
+                       "Successfully assigned %s to node %s", alias, nodeName)
 
                // before binding pod to node, first bind volumes to pod
                log.Log(log.ShimCacheTask).Debug("bind pod volumes",
-                       zap.String("podName", task.pod.Name),
-                       zap.String("podUID", string(task.pod.UID)))
-               if err := task.context.bindPodVolumes(task.pod); err != nil {
-                       log.Log(log.ShimCacheTask).Error("bind volumes to pod 
failed", zap.String("taskID", task.taskID), zap.Error(err))
-                       task.failWithEvent(fmt.Sprintf("bind volumes to pod 
failed, name: %s, %s", task.alias, err.Error()), "PodVolumesBindFailure")
+                       zap.String("podName", pod.Name),
+                       zap.String("podUID", string(pod.UID)))
+               if err := retry.OnError(bindPodBackoff, func(err error) bool {
+                       log.Log(log.ShimCacheTask).Error("bind volumes to pod 
failed, retrying",
+                               zap.String("taskID", task.taskID), 
zap.Error(err))
+                       return true
+               }, func() error {
+                       return task.context.bindPodVolumes(pod)
+               }); err != nil {
+                       log.Log(log.ShimCacheTask).Error("bind volumes to pod 
failed after retries",
+                               zap.String("taskID", task.taskID), 
zap.Error(err))
+                       task.rescheduleOnBindFailure(allocationKey, nodeName, 
"PodVolumesBindFailure",
+                               fmt.Sprintf("Failed to bind volumes for %s on 
node %s, it will be retried", alias, nodeName))
                        return
                }
                log.Log(log.ShimCacheTask).Debug("bind pod",
-                       zap.String("podName", task.pod.Name),
-                       zap.String("podUID", string(task.pod.UID)))
+                       zap.String("podName", pod.Name),
+                       zap.String("podUID", string(pod.UID)))
 
-               if err := 
task.context.apiProvider.GetAPIs().KubeClient.Bind(task.pod, task.nodeName); 
err != nil {
-                       log.Log(log.ShimCacheTask).Error("bind pod to node 
failed", zap.String("taskID", task.taskID), zap.Error(err))
-                       task.failWithEvent(fmt.Sprintf("bind pod to node 
failed, name: %s, %s", task.alias, err.Error()), "PodBindFailure")
+               if err := retry.OnError(bindPodBackoff, func(err error) bool {
+                       log.Log(log.ShimCacheTask).Error("bind pod to node 
failed, retrying",
+                               zap.String("taskID", task.taskID), 
zap.Error(err))
+                       return true
+               }, func() error {
+                       return 
task.context.apiProvider.GetAPIs().KubeClient.Bind(pod, nodeName)
+               }); err != nil {
+                       log.Log(log.ShimCacheTask).Error("bind pod to node 
failed after retries",
+                               zap.String("taskID", task.taskID), 
zap.Error(err))
+                       task.rescheduleOnBindFailure(allocationKey, nodeName, 
"PodBindFailure",
+                               fmt.Sprintf("Failed to bind %s to node %s, it 
will be retried", alias, nodeName))
                        return
                }
-               log.Log(log.ShimCacheTask).Info("successfully bound pod", 
zap.String("podName", task.pod.Name))
-               dispatcher.Dispatch(NewBindTaskEvent(task.applicationID, 
task.taskID))
-               events.GetRecorder().Eventf(task.pod.DeepCopy(), nil,
-                       v1.EventTypeNormal, "PodBindSuccessful", 
"PodBindSuccessful",
-                       "Pod %s is successfully bound to node %s", task.alias, 
task.nodeName)
+               log.Log(log.ShimCacheTask).Info("successfully bound pod", 
zap.String("podName", pod.Name))
 
+               task.lock.Lock()
                task.schedulingState = TaskSchedAllocated
-       }()
+               task.lock.Unlock()

Review Comment:
   This lock needs to be taken at the start of the go routine as we do not want 
anything to manipulate the task while running. It has to be a write lock also.
   The removal of the lock outside the routine is correct but we need the lock 
inside the go routine as the first thing. The go routine will be blocked until 
the `handle()` returns as expected



##########
pkg/cache/task.go:
##########
@@ -669,6 +709,22 @@ func (task *Task) 
rollbackOnAssumePodFailure(allocationKey, nodeID string) {
                zap.String("allocationKey", allocationKey))
 }
 
+// rescheduleOnBindFailure is called when volume or pod binding fails after 
all retries.
+// Move the task back to Scheduling before rolling back the allocation, so a
+// subsequent TaskAllocated event from the core (which requires Scheduling) is 
accepted.
+// Must be called without holding the task lock.
+func (task *Task) rescheduleOnBindFailure(allocationKey, nodeID, eventReason, 
eventMsg string) {
+       // Move the task back to Scheduling before releasing to the core, so 
the re-delivered
+       // allocation (valid only from the Scheduling state) is accepted by the 
state machine.
+       if err := task.handle(NewRescheduleTaskEvent(task.applicationID, 
task.taskID)); err != nil {

Review Comment:
   You must hold the task lock to do the (volume)binding. The calling go 
routine has the write lock. The handle will block as it needs a read lock.
   The pattern used in the k8shim is not good for this. You either need to 
dispatch an event or simply change the state back via 
`task.sm.SetState(states.Scheduling). Preference is a direct state setting as 
it has less overhead than an event dispatch.
   



##########
pkg/cache/task_state.go:
##########
@@ -47,10 +47,11 @@ const (
        TaskFail
        KillTask
        TaskKilled
+       TaskRescheduling

Review Comment:
   See above, direct state setting is preferred.



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