Unverified Commit e690b058 authored by Lemon-miaow's avatar Lemon-miaow
Browse files

fix(restore): wait for the tracking finalizer before recreating

Live verification of the previous commit showed the immediate retry STILL
stranded: deleting a finished Job leaves it terminating (job-tracking
finalizer), so the re-Create collided with the dying object and was
mapped to ErrAlreadyExists a second time. Poll until the name actually
frees (bounded, ~10s) and surface a 'retry shortly' error if a stuck
finalizer ever outlives the budget. Fake-client tests pin both the
replace-finished and coalesce-in-flight branches.
parent 90ccbfed
Loading
Loading
Loading
Loading
+37 −0
Changes for internal/restore/k8sjobs.go: 37 added lines, 0 removed lines.
Original line number Diff line number Diff line
@@ -2,6 +2,8 @@ package restore

import (
	"context"
	"fmt"
	"time"

	batchv1 "k8s.io/api/batch/v1"
	apierrors "k8s.io/apimachinery/pkg/api/errors"
@@ -71,9 +73,44 @@ func (k *K8sJobs) CreateRestoreJob(ctx context.Context, p JobParams) error {
	if deleteErr := k.c.Delete(ctx, &existing); deleteErr != nil && !apierrors.IsNotFound(deleteErr) {
		return deleteErr
	}
	// The API server keeps the object until its job-tracking finalizer has run,
	// so an immediate re-Create would collide again and swallow the retry a second
	// time (found live: the E2E retry still answered 202 while nothing ran). Wait
	// for the name to actually free, bounded, then replace.
	if err := k.waitForNameRelease(ctx, job.Namespace, job.Name); err != nil {
		return err
	}
	return k.recreate(ctx, job)
}

// waitForNameRelease polls until the named Job is gone or the wait budget is
// spent. The job controller releases the tracking finalizer within a second or
// two of the delete, so this normally returns on the first or second probe; the
// bound exists so a stuck finalizer surfaces as an error ("retry shortly")
// instead of another silent success.
func (k *K8sJobs) waitForNameRelease(ctx context.Context, namespace, name string) error {
	const (
		probeInterval = 500 * time.Millisecond
		maxProbes     = 20
	)
	for probe := 0; probe < maxProbes; probe++ {
		var probeJob batchv1.Job
		err := k.c.Get(ctx, types.NamespacedName{Namespace: namespace, Name: name}, &probeJob)
		if apierrors.IsNotFound(err) {
			return nil
		}
		if err != nil {
			return err
		}
		select {
		case <-ctx.Done():
			return ctx.Err()
		case <-time.After(probeInterval):
		}
	}
	return fmt.Errorf("restore: job %s/%s is still terminating after deletion; retry shortly", namespace, name)
}

// recreate retries Create once after a finished Job released the name. A
// collision that survives means a concurrent restore re-created first, so the
// idempotent answer applies again.
+72 −0
Changes for internal/restore/k8sjobs_test.go: 72 added lines, 0 removed lines.
Original line number Diff line number Diff line
package restore

import (
	"context"
	"errors"
	"testing"
	"time"

	batchv1 "k8s.io/api/batch/v1"
	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"k8s.io/apimachinery/pkg/runtime"
	"k8s.io/apimachinery/pkg/types"
	"sigs.k8s.io/controller-runtime/pkg/client/fake"
)

func testScheme(t *testing.T) *runtime.Scheme {
	t.Helper()
	scheme := runtime.NewScheme()
	if err := batchv1.AddToScheme(scheme); err != nil {
		t.Fatalf("scheme: %v", err)
	}
	return scheme
}

func testParams() JobParams {
	return JobParams{
		Server:    "survival",
		WorldPVC:  "world-survival-0",
		BackupPVC: "felis-backups",
		Namespace: "minecraft",
		Image:     "felis:test",
		BackupRef: "/backups/x.tar.gz",
	}
}

// A finished Job still holding the deterministic name must be replaced, not
// treated as an in-flight coalesce — otherwise the retry after a failed restore
// is answered 202 while nothing runs (found by an E2E audit).
func TestCreateRestoreJobReplacesFinishedJob(t *testing.T) {
	finished := &batchv1.Job{
		ObjectMeta: metav1.ObjectMeta{Name: "restore-survival", Namespace: "minecraft"},
		Status: batchv1.JobStatus{
			Failed:     1,
			Conditions: []batchv1.JobCondition{{Type: batchv1.JobFailed, Status: "True"}},
		},
	}
	c := fake.NewClientBuilder().WithScheme(testScheme(t)).WithObjects(finished).Build()

	if err := NewK8sJobs(c).CreateRestoreJob(context.Background(), testParams()); err != nil {
		t.Fatalf("CreateRestoreJob: %v", err)
	}
	var got batchv1.Job
	if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "restore-survival"}, &got); err != nil {
		t.Fatalf("replacement job missing: %v", err)
	}
	if restoreJobFinished(&got) {
		t.Errorf("replacement job is already finished: %+v", got.Status)
	}
}

// An in-flight Job keeps the idempotent coalesce: a duplicate enqueue during a
// running restore is absorbed, and the running Job is left untouched.
func TestCreateRestoreJobCoalescesInFlightJob(t *testing.T) {
	inFlight := &batchv1.Job{
		ObjectMeta: metav1.ObjectMeta{Name: "restore-survival", Namespace: "minecraft"},
		Status:     batchv1.JobStatus{Active: 1},
	}
	c := fake.NewClientBuilder().WithScheme(testScheme(t)).WithObjects(inFlight).Build()

	err := NewK8sJobs(c).CreateRestoreJob(context.Background(), testParams())
	if !errors.Is(err, ErrAlreadyExists) {
		t.Fatalf("CreateRestoreJob = %v, want ErrAlreadyExists", err)
	}
	var got batchv1.Job
	if err := c.Get(context.Background(), types.NamespacedName{Namespace: "minecraft", Name: "restore-survival"}, &got); err != nil {
		t.Fatalf("in-flight job should stay: %v", err)
	}
	if got.Status.Active != 1 {
		t.Errorf("in-flight job was disturbed: %+v", got.Status)
	}
}

// restoreJobFinished decides whether a name collision is a genuine in-flight
// coalesce (ErrAlreadyExists) or a finished Job whose deterministic name must be
// replaced so a retry enqueues for real. Regression: a FAILED restore used to