From b68aa1534446358bf21f69c315f93f513636f096 Mon Sep 17 00:00:00 2001 From: Lyndon-Li Date: Thu, 18 Jun 2026 16:41:04 +0800 Subject: [PATCH] wait PV detachment before deleting PV Signed-off-by: Lyndon-Li --- pkg/exposer/generic_restore.go | 7 ++++ pkg/util/kube/pvc_pv.go | 68 +++++++++++++++++++++++++++++----- 2 files changed, 65 insertions(+), 10 deletions(-) diff --git a/pkg/exposer/generic_restore.go b/pkg/exposer/generic_restore.go index f0ad76123..5d0f34d99 100644 --- a/pkg/exposer/generic_restore.go +++ b/pkg/exposer/generic_restore.go @@ -452,6 +452,13 @@ func (e *genericRestoreExposer) RebindVolume(ctx context.Context, ownerObject co curLog.WithField("restore PVC", restorePVCName).Info("Restore PVC is deleted") + err = kube.WaitVolumeDetached(ctx, e.kubeClient.StorageV1(), retained.Name, param.OperationTimeout) + if err != nil { + return errors.Wrapf(err, "error waiting for retained PV %s to detach", retained.Name) + } + + curLog.WithField("retained PV", retained.Name).Info("Retained PV is detached") + rebindPV, err = kube.RebindPV(ctx, e.kubeClient.CoreV1(), uuid.NewString(), retained, targetPVC, orgReclaim, param.TargetFSType) if err != nil { return errors.Wrapf(err, "error rebinding PV for target PVC %s", param.TargetPVCName) diff --git a/pkg/util/kube/pvc_pv.go b/pkg/util/kube/pvc_pv.go index 7dea36f08..182b18995 100644 --- a/pkg/util/kube/pvc_pv.go +++ b/pkg/util/kube/pvc_pv.go @@ -324,15 +324,10 @@ func RebindPV(ctx context.Context, pvGetter corev1client.CoreV1Interface, pvName maps.Copy(pvLabel, pvc.Spec.Selector.MatchLabels) } - pvAnnotations := make(map[string]string) - maps.Copy(pvAnnotations, source.Annotations) - delete(pvAnnotations, KubeAnnBoundByController) - pv := &corev1api.PersistentVolume{ ObjectMeta: metav1.ObjectMeta{ - Name: pvName, - Labels: pvLabel, - Annotations: pvAnnotations, + Name: pvName, + Labels: pvLabel, }, Spec: corev1api.PersistentVolumeSpec{ Capacity: source.Spec.Capacity, @@ -358,6 +353,10 @@ func RebindPV(ctx context.Context, pvGetter corev1client.CoreV1Interface, pvName func clonePVSource(source *corev1api.PersistentVolumeSource, newFSType string) corev1api.PersistentVolumeSource { newSource := source.DeepCopy() + if newSource.CSI != nil && newSource.CSI.VolumeAttributes != nil { + delete(newSource.CSI.VolumeAttributes, "storage.kubernetes.io/csiProvisionerIdentity") + } + if newFSType != "" { if newSource.CSI != nil { newSource.CSI.FSType = newFSType @@ -684,22 +683,71 @@ func GetPVAttachedNode(ctx context.Context, pv string, storageClient storagev1.S return "", nil } -func GetPVAttachedNodes(ctx context.Context, pv string, storageClient storagev1.StorageV1Interface) ([]string, error) { +func getPVAttachment(ctx context.Context, pv string, storageClient storagev1.StorageV1Interface) ([]*storagev1api.VolumeAttachment, error) { vaList, err := storageClient.VolumeAttachments().List(ctx, metav1.ListOptions{}) if err != nil { return nil, errors.Wrapf(err, "error listing volumeattachment") } - nodes := []string{} + attachments := []*storagev1api.VolumeAttachment{} for _, va := range vaList.Items { if va.Spec.Source.PersistentVolumeName != nil && *va.Spec.Source.PersistentVolumeName == pv { - nodes = append(nodes, va.Spec.NodeName) + attachments = append(attachments, &va) } } + return attachments, nil +} + +func GetPVAttachedNodes(ctx context.Context, pv string, storageClient storagev1.StorageV1Interface) ([]string, error) { + attachments, err := getPVAttachment(ctx, pv, storageClient) + if err != nil { + return nil, errors.Wrap(err, "error listing volumeattachment") + } + + nodes := []string{} + for _, attach := range attachments { + nodes = append(nodes, attach.Spec.NodeName) + } + return nodes, nil } +func WaitVolumeDetached(ctx context.Context, storageClient storagev1.StorageV1Interface, pv string, timeout time.Duration) error { + attachments, err := getPVAttachment(ctx, pv, storageClient) + if err != nil { + return errors.Wrap(err, "error listing volumeattachment") + } + + err = wait.PollUntilContextTimeout(ctx, waitInternal, timeout, true, func(ctx context.Context) (bool, error) { + left := []*storagev1api.VolumeAttachment{} + for _, attach := range attachments { + if _, err := storageClient.VolumeAttachments().Get(ctx, attach.Name, metav1.GetOptions{}); err == nil { + left = append(left, attach) + } else if !apierrors.IsNotFound(err) { + return false, err // Return the error if it's not a NotFound error + } + } + + if len(left) == 0 { + return true, nil + } + + attachments = left + + return false, nil + }) + + if err != nil { + if errors.Is(err, context.DeadlineExceeded) { + return errors.Errorf("timeout waiting for volume %s to be detached", pv) + } + return errors.Wrapf(err, "error waiting for volume %s to be detached", pv) + } + + return nil +} + func GetVolumeTopology(ctx context.Context, volumeClient corev1client.CoreV1Interface, storageClient storagev1.StorageV1Interface, pvName string, scName string) (*corev1api.NodeSelector, error) { if pvName == "" || scName == "" { return nil, errors.Errorf("invalid parameter, pv %s, sc %s", pvName, scName)