Skip to content
Closed
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,12 +28,20 @@ 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
// OperatorStale reports whether the training operator has shown no sign of life on a job for longer than timeout.
// The operator stamps StartTime once it creates the replica pods, after which the job is out of its hands. Before
// that, a gang-scheduled job whose PodGroup is still waiting for scheduler admission is held before pod creation
// and receives only LastReconcileTime, which the operator refreshes on every reconcile while it waits. A job with
// no StartTime whose most recent activity (creation or last reconcile) is older than timeout has been abandoned.
func OperatorStale(created meta_v1.Time, status kubeflowv1.JobStatus, timeout time.Duration, now time.Time) bool {
if status.StartTime != nil {
return false
}
lastSeen := created.Time
if status.LastReconcileTime != nil && status.LastReconcileTime.After(lastSeen) {
lastSeen = status.LastReconcileTime.Time
}
return lastSeen.Add(timeout).Before(now)
}

// ExtractCurrentCondition will return the first job condition for tensorflow/pytorch
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,16 +28,32 @@ func TestMain(m *testing.M) {
os.Exit(code)
}

func TestOperatorNeverReconciled(t *testing.T) {
now := meta_v1.Now()
func TestOperatorStale(t *testing.T) {
now := time.Now()
timeout := time.Minute
ts := func(ago time.Duration) *meta_v1.Time { return &meta_v1.Time{Time: now.Add(-ago)} }
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}))
// Never touched by the operator: stale only once the timeout has elapsed since creation.
assert.False(t, OperatorStale(*ts(time.Second), kubeflowv1.JobStatus{}, timeout, now))
assert.True(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{}, timeout, now))
assert.True(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{Conditions: []kubeflowv1.JobCondition{created}}, timeout, now))

// StartTime means pods exist; the operator's liveness no longer matters for this check.
assert.False(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{StartTime: ts(time.Hour)}, timeout, now))
assert.False(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{StartTime: ts(time.Hour), LastReconcileTime: ts(time.Hour)}, timeout, now))

// Gang-scheduled job waiting for PodGroup admission: pods (and StartTime) are withheld, but the operator keeps
// refreshing LastReconcileTime while it waits.
assert.False(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{LastReconcileTime: ts(time.Second)}, timeout, now))
assert.False(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{LastReconcileTime: ts(timeout - time.Second)}, timeout, now))

// Operator reconciled once and then went away: the old LastReconcileTime does not keep the job alive.
assert.True(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{LastReconcileTime: ts(timeout + time.Second)}, timeout, now))
assert.True(t, OperatorStale(*ts(time.Hour), kubeflowv1.JobStatus{LastReconcileTime: ts(30 * time.Minute)}, timeout, now))

// A LastReconcileTime older than creation (clock skew) never shortens the grace period.
assert.False(t, OperatorStale(*ts(time.Second), kubeflowv1.JobStatus{LastReconcileTime: ts(time.Hour)}, timeout, now))
}

func TestExtractCurrentCondition(t *testing.T) {
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 && common.OperatorNeverReconciled(app.Status) && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
if !isSuspended && common.OperatorStale(app.CreationTimestamp, app.Status, common.GetConfig().Timeout.Duration, 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
15 changes: 13 additions & 2 deletions flyteplugins/go/tasks/plugins/k8s/kfoperators/mpi/mpi_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -594,15 +594,26 @@ func TestGetTaskPhase(t *testing.T) {
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
// admitted: StartTime stays nil past the timeout while a fresh 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)}
mpiJobGangWait.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-time.Second)}
taskPhase, err = mpiResourceHandler.GetTaskPhase(ctx, taskCtx, mpiJobGangWait)
assert.NoError(t, err)
assert.Equal(t, pluginsCore.PhaseQueued, taskPhase.Phase())

// Training operator reconciled the job once and then stopped: a LastReconcileTime older than the timeout
// does not count as liveness
mpiJobOperatorGone := dummyMPIJobResourceCreator(kubeflowv1.JobCreated)
mpiJobOperatorGone.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
mpiJobOperatorGone.Status.StartTime = nil
mpiJobOperatorGone.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-30 * time.Minute)}
taskPhase, err = mpiResourceHandler.GetTaskPhase(ctx, taskCtx, mpiJobOperatorGone)
assert.Error(t, err)
assert.Contains(t, err.Error(), "kubeflow operator hasn't updated")
assert.Equal(t, pluginsCore.PhaseInfoUndefined, taskPhase)

// 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 && common.OperatorNeverReconciled(app.Status) && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
if !isSuspended && common.OperatorStale(app.CreationTimestamp, app.Status, common.GetConfig().Timeout.Duration, 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 @@ -715,15 +715,26 @@ func TestGetTaskPhase(t *testing.T) {
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
// admitted: StartTime stays nil past the timeout while a fresh 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)}
pytorchJobGangWait.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-time.Second)}
taskPhase, err = pytorchResourceHandler.GetTaskPhase(ctx, taskCtx, pytorchJobGangWait)
assert.NoError(t, err)
assert.Equal(t, pluginsCore.PhaseQueued, taskPhase.Phase())

// Training operator reconciled the job once and then stopped: a LastReconcileTime older than the timeout
// does not count as liveness
pytorchJobOperatorGone := dummyPytorchJobResourceCreator(kubeflowv1.JobCreated)
pytorchJobOperatorGone.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
pytorchJobOperatorGone.Status.StartTime = nil
pytorchJobOperatorGone.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-30 * time.Minute)}
taskPhase, err = pytorchResourceHandler.GetTaskPhase(ctx, taskCtx, pytorchJobOperatorGone)
assert.Error(t, err)
assert.Contains(t, err.Error(), "kubeflow operator hasn't updated")
assert.Equal(t, pluginsCore.PhaseInfoUndefined, taskPhase)

// 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 && common.OperatorNeverReconciled(app.Status) && app.CreationTimestamp.Add(common.GetConfig().Timeout.Duration).Before(time.Now()) {
if !isSuspended && common.OperatorStale(app.CreationTimestamp, app.Status, common.GetConfig().Timeout.Duration, 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 @@ -640,15 +640,26 @@ func TestGetTaskPhase(t *testing.T) {
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
// admitted: StartTime stays nil past the timeout while a fresh 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)}
tfJobGangWait.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-time.Second)}
taskPhase, err = tensorflowResourceHandler.GetTaskPhase(ctx, taskCtx, tfJobGangWait)
assert.NoError(t, err)
assert.Equal(t, pluginsCore.PhaseQueued, taskPhase.Phase())

// Training operator reconciled the job once and then stopped: a LastReconcileTime older than the timeout
// does not count as liveness
tfJobOperatorGone := dummyTensorFlowJobResourceCreator(kubeflowv1.JobCreated)
tfJobOperatorGone.CreationTimestamp = v1.Time{Time: time.Now().Add(-time.Hour)}
tfJobOperatorGone.Status.StartTime = nil
tfJobOperatorGone.Status.LastReconcileTime = &v1.Time{Time: time.Now().Add(-30 * time.Minute)}
taskPhase, err = tensorflowResourceHandler.GetTaskPhase(ctx, taskCtx, tfJobOperatorGone)
assert.Error(t, err)
assert.Contains(t, err.Error(), "kubeflow operator hasn't updated")
assert.Equal(t, pluginsCore.PhaseInfoUndefined, taskPhase)

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