From 073edc81cb664eef021657cdbab241fe4d7cdb6c Mon Sep 17 00:00:00 2001 From: "shenmu.wy" Date: Wed, 16 Sep 2026 15:35:29 +0800 Subject: [PATCH] feat: support the EtcdMember Replacing workflow Signed-off-by: shenmu.wy --- internal/controller/etcdmember.go | 59 ++++++++++++++-- internal/controller/etcdmember_test.go | 78 +++++++++++++++++++++ test/e2e/e2e_test.go | 94 ++++++++++++++++++++++++++ 3 files changed, 225 insertions(+), 6 deletions(-) diff --git a/internal/controller/etcdmember.go b/internal/controller/etcdmember.go index 31fcaf85..82aed044 100644 --- a/internal/controller/etcdmember.go +++ b/internal/controller/etcdmember.go @@ -78,12 +78,56 @@ func (r *EtcdClusterReconciler) reconcileEtcdMember( // releasing the finalizer so Kubernetes can finish the deletion. case ecv1alpha1.EtcdMemberTerminating: return r.cleanupEtcdMember(ctx, state, member) - // placeholder for Recreating and Replacing + case ecv1alpha1.EtcdMemberReplacing: + return r.replaceEtcdMember(ctx, state, member) + // placeholder for Recreating default: return ctrl.Result{}, nil } } +// markMemberReplacing persists the replacement intent before cleanup starts. +// Resetting the count gives the replacement a fresh provisioning retry budget. +// +//nolint:unused // The shared recovery ladder will call this when #472 is implemented. +func (r *EtcdClusterReconciler) markMemberReplacing(ctx context.Context, member *ecv1alpha1.EtcdMember) error { + return r.updateEtcdMemberStatus(ctx, member, func(status *ecv1alpha1.EtcdMemberStatus) { + status.Phase = ecv1alpha1.EtcdMemberReplacing + status.RecreateCount = 0 + }) +} + +// replaceEtcdMember cleans up the old member and waits for its Pod and PVC +// to disappear before rejoining through the existing provisioning workflow. +func (r *EtcdClusterReconciler) replaceEtcdMember(ctx context.Context, state *reconcileState, member *ecv1alpha1.EtcdMember) (ctrl.Result, error) { + if _, err := r.cleanupEtcdMember(ctx, state, member); err != nil { + return ctrl.Result{}, err + } + + // Both Pod and PVC must be deleted before provisioning new ones. + resources := []client.Object{ + &corev1.Pod{ObjectMeta: metav1.ObjectMeta{Name: member.Name, Namespace: member.Namespace}}, + &corev1.PersistentVolumeClaim{ObjectMeta: metav1.ObjectMeta{Name: pvcNameForMember(member.Name), Namespace: member.Namespace}}, + } + for _, resource := range resources { + if err := r.Get(ctx, client.ObjectKeyFromObject(resource), resource); err != nil { + if apierrors.IsNotFound(err) { + continue + } + return ctrl.Result{}, fmt.Errorf("checking resource %q for replacing EtcdMember %q: %w", resource.GetName(), member.Name, err) + } + return ctrl.Result{RequeueAfter: requeueDuration}, nil + } + + if err := r.updateEtcdMemberStatus(ctx, member, func(status *ecv1alpha1.EtcdMemberStatus) { + status.Phase = ecv1alpha1.EtcdMemberProvisioning + }); err != nil { + return ctrl.Result{}, err + } + // Joining must use a fresh membership snapshot on the next reconcile. + return ctrl.Result{RequeueAfter: requeueDuration}, nil +} + // markMemberTerminating persists Phase=Terminating on a member whose // deletion has started, if not already written. Idempotent across repeated // dispatch attempts. @@ -96,9 +140,10 @@ func (r *EtcdClusterReconciler) markMemberTerminating(ctx context.Context, membe }) } -// cleanupEtcdMember is the §4.6 Terminating leave for one member: remove it -// from etcd's live membership, disarm its alarms, delete its owned Pod and PVC, -// and finally release memberCleanupFinalizer so Kubernetes can finish the deletion. +// cleanupEtcdMember is the §4.6 leave for one member: remove it from etcd's +// live membership, disarm its alarms, and delete its owned Pod and PVC. +// Only the Terminating case removes the memberCleanupFinalizer +// so Kubernetes can finish the deletion. // Membership and resource cleanup are no-ops once done, so re-entering after // an interruption (operator restart, transient etcd error) resumes harmlessly. // Alarm cleanup currently requires the pre-removal membership snapshot; @@ -137,8 +182,10 @@ func (r *EtcdClusterReconciler) cleanupEtcdMember(ctx context.Context, s *reconc return ctrl.Result{}, err } - if err := r.clearMemberFinalizer(ctx, member); err != nil { - return ctrl.Result{}, err + if member.Status.Phase == ecv1alpha1.EtcdMemberTerminating { + if err := r.clearMemberFinalizer(ctx, member); err != nil { + return ctrl.Result{}, err + } } return ctrl.Result{RequeueAfter: requeueDuration}, nil } diff --git a/internal/controller/etcdmember_test.go b/internal/controller/etcdmember_test.go index 59eec34c..b25be1e2 100644 --- a/internal/controller/etcdmember_test.go +++ b/internal/controller/etcdmember_test.go @@ -485,6 +485,7 @@ func TestCleanupEtcdMember(t *testing.T) { member := leaveTestMember(2) member.DeletionTimestamp = &now member.Finalizers = []string{memberCleanupFinalizer} + member.Status.Phase = ecv1alpha1.EtcdMemberTerminating pod := leaveTestPod(2) pvc := &corev1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{Name: pvcNameForMember(pod.Name), Namespace: ec.Namespace}, @@ -520,6 +521,83 @@ func TestCleanupEtcdMember(t *testing.T) { assert.Error(t, fakeClient.Get(ctx, types.NamespacedName{Namespace: ec.Namespace, Name: member.Name}, &ecv1alpha1.EtcdMember{})) } +// TestReplaceEtcdMember verifies that replacement waits for complete Pod/PVC +// deletion, retains the member finalizer, and resumes directly in Provisioning. +func TestReplaceEtcdMember(t *testing.T) { + ctx := t.Context() + scheme := leaveTestScheme(t) + ec := leaveTestCluster() + member := leaveTestMember(2) + member.Status.Phase = ecv1alpha1.EtcdMemberReplacing + member.Finalizers = []string{memberCleanupFinalizer} + pod := leaveTestPod(2) + pod.Finalizers = []string{"test.etcd.io/hold"} + pvc := &corev1.PersistentVolumeClaim{ObjectMeta: metav1.ObjectMeta{ + Name: pvcNameForMember(pod.Name), Namespace: ec.Namespace, + Finalizers: []string{"test.etcd.io/hold"}, + }} + fakeClient := fake.NewClientBuilder().WithScheme(scheme). + WithStatusSubresource(&ecv1alpha1.EtcdMember{}). + WithObjects(ec, member, pod, pvc).Build() + r := &EtcdClusterReconciler{Client: fakeClient, Scheme: scheme} + state := &reconcileState{ + cluster: ec, pods: []*corev1.Pod{pod}, + // Resume after the old etcd identity has been removed. + memberListResp: &clientv3.MemberListResponse{}, + } + assertReplacing := func() { + t.Helper() + require.NoError(t, fakeClient.Get(ctx, client.ObjectKeyFromObject(member), member)) + assert.Equal(t, ecv1alpha1.EtcdMemberReplacing, member.Status.Phase) + assert.Equal(t, []string{memberCleanupFinalizer}, member.Finalizers) + } + + // Repeated passes must wait for the Pod and leave the PVC untouched. + for range 2 { + res, err := r.replaceEtcdMember(ctx, state, member) + require.NoError(t, err) + assert.Equal(t, ctrl.Result{RequeueAfter: requeueDuration}, res) + require.NoError(t, fakeClient.Get(ctx, client.ObjectKeyFromObject(pod), pod)) + assert.NotNil(t, pod.DeletionTimestamp) + require.NoError(t, fakeClient.Get(ctx, client.ObjectKeyFromObject(pvc), pvc)) + assert.Nil(t, pvc.DeletionTimestamp) + assertReplacing() + } + pod.Finalizers = nil + require.NoError(t, fakeClient.Update(ctx, pod)) + state.pods = nil + + // Once the Pod is gone, a terminating PVC still blocks provisioning. + for range 2 { + res, err := r.replaceEtcdMember(ctx, state, member) + require.NoError(t, err) + assert.Equal(t, ctrl.Result{RequeueAfter: requeueDuration}, res) + require.NoError(t, fakeClient.Get(ctx, client.ObjectKeyFromObject(pvc), pvc)) + assert.NotNil(t, pvc.DeletionTimestamp) + assertReplacing() + } + pvc.Finalizers = nil + require.NoError(t, fakeClient.Update(ctx, pvc)) + + // Exercise phase dispatch after cleanup completes. It must persist + // Provisioning and requeue without creating resources in this pass. + res, err := r.reconcileEtcdMember(ctx, state, member) + require.NoError(t, err) + assert.Equal(t, ctrl.Result{RequeueAfter: requeueDuration}, res) + require.NoError(t, fakeClient.Get(ctx, client.ObjectKeyFromObject(member), member)) + assert.Equal(t, ecv1alpha1.EtcdMemberProvisioning, member.Status.Phase) + assert.Equal(t, []string{memberCleanupFinalizer}, member.Finalizers) + assert.Nil(t, member.DeletionTimestamp) + assert.Zero(t, member.Status.RecreateCount) + assert.Equal(t, leaveTestMember(2).Spec, member.Spec) + pods := &corev1.PodList{} + require.NoError(t, fakeClient.List(ctx, pods)) + assert.Empty(t, pods.Items) + pvcs := &corev1.PersistentVolumeClaimList{} + require.NoError(t, fakeClient.List(ctx, pvcs)) + assert.Empty(t, pvcs.Items) +} + // TestMarkMemberTerminating verifies the Phase write before the leave runs, // and that a repeated call is a no-op. func TestMarkMemberTerminating(t *testing.T) { diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go index 0029415e..d2a9a3d7 100644 --- a/test/e2e/e2e_test.go +++ b/test/e2e/e2e_test.go @@ -29,6 +29,7 @@ import ( "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/util/retry" "sigs.k8s.io/e2e-framework/klient/k8s/resources" "sigs.k8s.io/e2e-framework/klient/wait" "sigs.k8s.io/e2e-framework/klient/wait/conditions" @@ -391,6 +392,99 @@ func TestEtcdMemberDeletionRecreatesResources(t *testing.T) { _ = testEnv.Test(t, feature.Feature()) } +// TestEtcdMemberReplacing sets ordinal 2 to Phase Replacing, +// preserving its EtcdMember while recreating its Pod and PVC. +func TestEtcdMemberReplacing(t *testing.T) { + feature := features.New("member-replacing") + clusterName := "etcd-member-replacing" + memberName := clusterName + "-2" + pvcName := "etcd-data-" + memberName + + feature.Setup(func(ctx context.Context, t *testing.T, c *envconf.Config) context.Context { + createEtcdClusterWithPVC(ctx, t, c, clusterName, 3) + if err := waitForAllEtcdMemberReady(t, c, etcdClusterRef(clusterName, 3)); err != nil { + t.Fatalf("cluster %s did not become Ready before replacement: %v", clusterName, err) + } + return ctx + }) + + feature.Assess("all members recover to Ready with a new Pod and PVC for ordinal 2", + func(ctx context.Context, t *testing.T, c *envconf.Config) context.Context { + before := takeResourceIDSnapshot(t, c, memberName, pvcName) + if before.podUID == "" || before.pvcUID == "" { + t.Fatal("expected non-empty Pod and PVC UIDs before replacement") + } + t.Logf("before replacement: Pod UID=%s, PVC UID=%s", before.podUID, before.pvcUID) + + // Persist replacement intent directly; the status update also + // enqueues the owning cluster through its EtcdMember watch. + if err := retry.RetryOnConflict(retry.DefaultRetry, func() error { + var target ecv1alpha1.EtcdMember + if err := c.Client().Resources().Get(ctx, memberName, namespace, &target); err != nil { + return err + } + target.Status.Phase = ecv1alpha1.EtcdMemberReplacing + target.Status.RecreateCount = 0 + return c.Client().Resources().UpdateStatus(ctx, &target) + }); err != nil { + t.Fatalf("failed to mark member %s Replacing: %v", memberName, err) + } + + // Wait for changed identities first: the original Ready state + // must not satisfy the post-replacement readiness assertion. + if err := wait.For(func(ctx context.Context) (bool, error) { + var pod corev1.Pod + if err := c.Client().Resources().Get(ctx, memberName, namespace, &pod); err != nil { + if errors.IsNotFound(err) { + return false, nil + } + return false, err + } + var pvc corev1.PersistentVolumeClaim + if err := c.Client().Resources().Get(ctx, pvcName, namespace, &pvc); err != nil { + if errors.IsNotFound(err) { + return false, nil + } + return false, err + } + var member ecv1alpha1.EtcdMember + if err := c.Client().Resources().Get(ctx, memberName, namespace, &member); err != nil { + return false, err + } + // Ready and the observed etcd ID are separate status writes. + // Wait for observation to catch up before comparing snapshots. + return pod.DeletionTimestamp == nil && pvc.DeletionTimestamp == nil && + pod.UID != before.podUID && pvc.UID != before.pvcUID && + member.Status.MemberID != "" && member.Status.MemberID != before.memberID, nil + }, wait.WithContext(ctx), wait.WithTimeout(5*time.Minute), wait.WithInterval(2*time.Second)); err != nil { + t.Fatalf("member %s did not get a new Pod, PVC and etcd identity: %v", memberName, err) + } + if err := waitForAllEtcdMemberReady(t, c, etcdClusterRef(clusterName, 3)); err != nil { + t.Fatalf("cluster %s did not recover all three members to Ready: %v", clusterName, err) + } + + after := takeResourceIDSnapshot(t, c, memberName, pvcName) + t.Logf("after replacement: Pod UID=%s, PVC UID=%s", after.podUID, after.pvcUID) + if after.podUID == before.podUID || after.pvcUID == before.pvcUID { + t.Error("replacement must change both Pod and PVC UIDs") + } + if after.memberUID != before.memberUID { + t.Error("replacement must retain the original EtcdMember object") + } + if after.memberID == before.memberID { + t.Error("replacement must rejoin with a new etcd member ID") + } + return ctx + }) + + feature.Teardown(func(ctx context.Context, t *testing.T, c *envconf.Config) context.Context { + cleanupEtcdCluster(ctx, t, c, clusterName) + return ctx + }) + + _ = testEnv.Test(t, feature.Feature()) +} + // TestManualScaleInWithoutEndpoints verifies the issue #463 acceptance case: // a deadlocked cluster can be manually scaled in even when none of its etcd // processes is reachable. Per the recovery recipe, the user shrinks