pbacsko commented on code in PR #1052:
URL: https://github.com/apache/yunikorn-k8shim/pull/1052#discussion_r3713630094
##########
pkg/cache/task_test.go:
##########
@@ -930,3 +934,94 @@ func TestCheckPodMetadataBeforeScheduling(t *testing.T) {
})
}
}
+
+// newRollbackTask creates a task in Scheduling state with allocationKey and
nodeName set,
+// ready to exercise rollbackOnAssumePodFailure.
+func newRollbackTask(ctx *Context, allocationKey, nodeID string) *Task {
+ app := NewApplication(appID1, queueNameA, testUser, testGroups,
map[string]string{},
+ ctx.apiProvider.GetAPIs().SchedulerAPI)
+ pod := &v1.Pod{
+ TypeMeta: metav1.TypeMeta{Kind: "Pod", APIVersion: "v1"},
+ ObjectMeta: metav1.ObjectMeta{Name: "rollback-pod", UID:
"rollback-uid"},
+ }
+ task := NewTask(allocationKey, app, ctx, pod)
+ task.sm.SetState(TaskStates().Scheduling)
+ task.lock.Lock()
+ task.allocationKey = allocationKey
+ task.nodeName = nodeID
+ task.lock.Unlock()
Review Comment:
This object only exists here at this point, no need to lock-unlock.
##########
pkg/cache/task.go:
##########
@@ -620,6 +620,58 @@ func (task *Task) failWithEvent(errorMessage, actionReason
string) {
dispatcher.Dispatch(NewFailTaskEvent(task.applicationID, task.taskID,
errorMessage))
}
+// rollbackOnAssumePodFailure is called when AssumePod fails after all retries.
+// It resets task 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) rollbackOnAssumePodFailure(allocationKey, nodeID string) {
+ // Read fields needed for event posting and release request under read
lock.
+ task.lock.RLock()
+ podCopy := task.pod.DeepCopy()
+ alias := task.alias
+ appID := task.applicationID
+ partition := task.application.partition
+ task.lock.RUnlock()
Review Comment:
I think it is better from readability if you merge L629-634 with L642-L645
under a single write lock.
##########
pkg/cache/context.go:
##########
@@ -897,6 +897,37 @@ func (ctx *Context) ForgetPod(name string) {
log.Log(log.ShimContext).Debug("unable to forget pod: not found in
cache", zap.String("pod", name))
}
+// RevertPodVolumeAssumptions undoes any PV/PVC assumptions made by the volume
binder
+// for the given pod on the given node. This is idempotent and safe to call
even if
+// AssumePodVolumes was never called or already reverted internally.
+func (ctx *Context) RevertPodVolumeAssumptions(podName, nodeID string) {
+ ctx.lock.Lock()
+ defer ctx.lock.Unlock()
+ pod := ctx.schedulerCache.GetPod(podName)
+ if pod == nil {
+ return
+ }
+ node := ctx.schedulerCache.GetNode(nodeID)
+ if node == nil {
+ return
+ }
+ podVolumeClaims, err :=
ctx.apiProvider.GetAPIs().VolumeBinder.GetPodVolumeClaims(ctx.klogger, pod)
+ if err != nil {
+ log.Log(log.ShimContext).Debug("RevertPodVolumeAssumptions:
failed to get pod volume claims",
Review Comment:
we better log it with an appropriate log level
##########
pkg/cache/context.go:
##########
@@ -897,6 +897,37 @@ func (ctx *Context) ForgetPod(name string) {
log.Log(log.ShimContext).Debug("unable to forget pod: not found in
cache", zap.String("pod", name))
}
+// RevertPodVolumeAssumptions undoes any PV/PVC assumptions made by the volume
binder
+// for the given pod on the given node. This is idempotent and safe to call
even if
+// AssumePodVolumes was never called or already reverted internally.
+func (ctx *Context) RevertPodVolumeAssumptions(podName, nodeID string) {
+ ctx.lock.Lock()
+ defer ctx.lock.Unlock()
+ pod := ctx.schedulerCache.GetPod(podName)
+ if pod == nil {
+ return
+ }
+ node := ctx.schedulerCache.GetNode(nodeID)
+ if node == nil {
+ return
+ }
+ podVolumeClaims, err :=
ctx.apiProvider.GetAPIs().VolumeBinder.GetPodVolumeClaims(ctx.klogger, pod)
+ if err != nil {
+ log.Log(log.ShimContext).Debug("RevertPodVolumeAssumptions:
failed to get pod volume claims",
+ zap.String("pod", podName), zap.Error(err))
+ return
+ }
+ podVolumes, _, err :=
ctx.apiProvider.GetAPIs().VolumeBinder.FindPodVolumes(ctx.klogger, pod,
podVolumeClaims, node.Node())
+ if err != nil || podVolumes == nil {
+ log.Log(log.ShimContext).Debug("RevertPodVolumeAssumptions:
failed to find pod volumes",
Review Comment:
we better log it with an appropriate log level
--
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]