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
59 changes: 53 additions & 6 deletions internal/controller/etcdmember.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment on lines +93 to +96
})
}

// 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.
Expand All @@ -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;
Expand Down Expand Up @@ -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
}
Expand Down
78 changes: 78 additions & 0 deletions internal/controller/etcdmember_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down Expand Up @@ -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) {
Expand Down
94 changes: 94 additions & 0 deletions test/e2e/e2e_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down