Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,14 @@ const (
PytorchTaskType = "pytorch"
)

// OperatorNeverReconciled reports whether the training operator has yet to write a timestamp to the job's status.
// The operator only stamps StartTime once it creates the replica pods; a gang-scheduled job whose PodGroup is
// still waiting for scheduler admission is held before pod creation and receives only LastReconcileTime. Either
// field proves the operator saw the job, so only a status with neither means the operator is absent.
func OperatorNeverReconciled(status kubeflowv1.JobStatus) bool {
return status.StartTime == nil && status.LastReconcileTime == nil
}

// ExtractCurrentCondition will return the first job condition for tensorflow/pytorch
func ExtractCurrentCondition(jobConditions []kubeflowv1.JobCondition) (kubeflowv1.JobCondition, error) {
if jobConditions != nil {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,18 @@ func TestMain(m *testing.M) {
os.Exit(code)
}

func TestOperatorNeverReconciled(t *testing.T) {
now := meta_v1.Now()
created := kubeflowv1.JobCondition{Type: kubeflowv1.JobCreated, Status: corev1.ConditionTrue}

assert.True(t, OperatorNeverReconciled(kubeflowv1.JobStatus{}))
assert.True(t, OperatorNeverReconciled(kubeflowv1.JobStatus{Conditions: []kubeflowv1.JobCondition{created}}))
assert.False(t, OperatorNeverReconciled(kubeflowv1.JobStatus{StartTime: &now}))
// Gang-scheduled job waiting for PodGroup admission: pods (and StartTime) are withheld, only LastReconcileTime lands.
assert.False(t, OperatorNeverReconciled(kubeflowv1.JobStatus{LastReconcileTime: &now}))
assert.False(t, OperatorNeverReconciled(kubeflowv1.JobStatus{StartTime: &now, LastReconcileTime: &now}))
}

func TestExtractCurrentCondition(t *testing.T) {
jobCreated := kubeflowv1.JobCondition{
Type: kubeflowv1.JobCreated,
Expand Down
2 changes: 1 addition & 1 deletion flyteplugins/go/tasks/plugins/k8s/kfoperators/mpi/mpi.go
Original file line number Diff line number Diff line change
Expand Up @@ -176,7 +176,7 @@ func (mpiOperatorResourceHandler) GetTaskPhase(ctx context.Context, pluginContex
}

isSuspended := app.Spec.RunPolicy.Suspend != nil && *app.Spec.RunPolicy.Suspend
if !isSuspended && app.Status.StartTime == nil && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
if !isSuspended && common.OperatorNeverReconciled(app.Status) && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
return pluginsCore.PhaseInfoUndefined, fmt.Errorf("kubeflow operator hasn't updated the mpi custom resource since creation time %v", app.CreationTimestamp)
}
currentCondition, err := common.ExtractCurrentCondition(app.Status.Conditions)
Expand Down
10 changes: 10 additions & 0 deletions flyteplugins/go/tasks/plugins/k8s/kfoperators/mpi/mpi_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -593,6 +593,16 @@ func TestGetTaskPhase(t *testing.T) {
assert.Contains(t, err.Error(), "kubeflow operator hasn't updated")
assert.Equal(t, pluginsCore.PhaseInfoUndefined, taskPhase)

// Training operator reconciled the job but is holding pod creation until the gang's PodGroup is
// admitted: StartTime stays nil past the timeout while LastReconcileTime proves the operator is alive
mpiJobGangWait := dummyMPIJobResourceCreator(kubeflowv1.JobCreated)
mpiJobGangWait.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
mpiJobGangWait.Status.StartTime = nil
mpiJobGangWait.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-time.Minute)}
taskPhase, err = mpiResourceHandler.GetTaskPhase(ctx, taskCtx, mpiJobGangWait)
assert.NoError(t, err)
assert.Equal(t, pluginsCore.PhaseQueued, taskPhase.Phase())

// Training operator did not modify the job because it is suspended
mpiJobSuspended := dummyMPIJobResourceCreator(kubeflowv1.JobCreated)
mpiJobSuspended.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,7 @@ func (pytorchOperatorResourceHandler) GetTaskPhase(ctx context.Context, pluginCo
}

isSuspended := app.Spec.RunPolicy.Suspend != nil && *app.Spec.RunPolicy.Suspend
if !isSuspended && app.Status.StartTime == nil && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
if !isSuspended && common.OperatorNeverReconciled(app.Status) && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
return pluginsCore.PhaseInfoUndefined, fmt.Errorf("kubeflow operator hasn't updated the pytorch custom resource since creation time %v", app.CreationTimestamp)
}
currentCondition, err := common.ExtractCurrentCondition(app.Status.Conditions)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -714,6 +714,16 @@ func TestGetTaskPhase(t *testing.T) {
assert.Contains(t, err.Error(), "kubeflow operator hasn't updated")
assert.Equal(t, pluginsCore.PhaseInfoUndefined, taskPhase)

// Training operator reconciled the job but is holding pod creation until the gang's PodGroup is
// admitted: StartTime stays nil past the timeout while LastReconcileTime proves the operator is alive
pytorchJobGangWait := dummyPytorchJobResourceCreator(kubeflowv1.JobCreated)
pytorchJobGangWait.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
pytorchJobGangWait.Status.StartTime = nil
pytorchJobGangWait.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-time.Minute)}
taskPhase, err = pytorchResourceHandler.GetTaskPhase(ctx, taskCtx, pytorchJobGangWait)
assert.NoError(t, err)
assert.Equal(t, pluginsCore.PhaseQueued, taskPhase.Phase())

// Training operator did not modify the job because it is suspended
pytorchJobSuspended := dummyPytorchJobResourceCreator(kubeflowv1.JobCreated)
pytorchJobSuspended.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,7 @@ func (tensorflowOperatorResourceHandler) GetTaskPhase(ctx context.Context, plugi
}

isSuspended := app.Spec.RunPolicy.Suspend != nil && *app.Spec.RunPolicy.Suspend
if !isSuspended && app.Status.StartTime == nil && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
if !isSuspended && common.OperatorNeverReconciled(app.Status) && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
return pluginsCore.PhaseInfoUndefined, fmt.Errorf("kubeflow operator hasn't updated the tensorflow custom resource since creation time %v", app.CreationTimestamp)
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -639,6 +639,16 @@ func TestGetTaskPhase(t *testing.T) {
assert.Contains(t, err.Error(), "kubeflow operator hasn't updated")
assert.Equal(t, pluginsCore.PhaseInfoUndefined, taskPhase)

// Training operator reconciled the job but is holding pod creation until the gang's PodGroup is
// admitted: StartTime stays nil past the timeout while LastReconcileTime proves the operator is alive
tfJobGangWait := dummyTensorFlowJobResourceCreator(kubeflowv1.JobCreated)
tfJobGangWait.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
tfJobGangWait.Status.StartTime = nil
tfJobGangWait.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-time.Minute)}
taskPhase, err = tensorflowResourceHandler.GetTaskPhase(ctx, taskCtx, tfJobGangWait)
assert.NoError(t, err)
assert.Equal(t, pluginsCore.PhaseQueued, taskPhase.Phase())

// Training operator did not modify the job because it is suspended
tfJobSuspended := dummyTensorFlowJobResourceCreator(kubeflowv1.JobCreated)
tfJobSuspended.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
Expand Down
Loading