diff --git a/internal/runtime/kubernetes/procedure_attempts.go b/internal/runtime/kubernetes/procedure_attempts.go index ae97b4c..a7d1ca5 100644 --- a/internal/runtime/kubernetes/procedure_attempts.go +++ b/internal/runtime/kubernetes/procedure_attempts.go @@ -181,7 +181,7 @@ func procedureJobAttempt(job *batchv1.Job, baseName string) int { } func (b *Backend) createFreshJob(ctx context.Context, job *batchv1.Job) (*batchv1.Job, error) { - propagation := metav1.DeletePropagationBackground + propagation := metav1.DeletePropagationForeground deleteCtx, cancelDelete := context.WithTimeout(ctx, 30*time.Second) defer cancelDelete() existing, err := b.client.BatchV1().Jobs(job.Namespace).Get(deleteCtx, job.Name, metav1.GetOptions{}) @@ -189,6 +189,10 @@ func (b *Backend) createFreshJob(ctx context.Context, job *batchv1.Job) (*batchv logger.Log().Error("Failed to check Kubernetes job before create", zap.String("namespace", job.Namespace), zap.String("job", job.Name), zap.Error(err)) return nil, err } + if err == nil && existing.Status.Succeeded == 0 && !kubernetesJobFailed(existing) { + logger.Log().Info("Reusing active Kubernetes job", zap.String("namespace", job.Namespace), zap.String("job", job.Name)) + return existing, nil + } if existing != nil && kubernetesJobFailed(existing) { original := job.Name job = job.DeepCopy() diff --git a/internal/runtime/kubernetes/resources_test.go b/internal/runtime/kubernetes/resources_test.go index a958466..baaeef1 100644 --- a/internal/runtime/kubernetes/resources_test.go +++ b/internal/runtime/kubernetes/resources_test.go @@ -535,6 +535,31 @@ func TestCreateFreshJobKeepsFailedJob(t *testing.T) { } } +func TestCreateFreshJobReusesPendingJob(t *testing.T) { + existing := &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "pending", Namespace: "druid", UID: "existing"}, + } + client := fake.NewSimpleClientset(existing) + backend := NewWithClient(Config{Namespace: "druid"}, client) + + created, err := backend.createFreshJob(context.Background(), &batchv1.Job{ + ObjectMeta: metav1.ObjectMeta{Name: "pending", Namespace: "druid"}, + }) + if err != nil { + t.Fatal(err) + } + if created.UID != existing.UID { + t.Fatalf("job UID = %q, want existing UID %q", created.UID, existing.UID) + } + jobs, err := client.BatchV1().Jobs("druid").List(context.Background(), metav1.ListOptions{}) + if err != nil { + t.Fatal(err) + } + if len(jobs.Items) != 1 { + t.Fatalf("jobs = %d, want 1", len(jobs.Items)) + } +} + func TestCreateOrReuseProcedureJobRetainsFailedBaseAndCreatesRetry(t *testing.T) { root := ref("druid", dataPVCName("deployment-123")) base := procedureResourceName(root, "start", 1)