Skip to content
Closed
Show file tree
Hide file tree
Changes from 4 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
13 changes: 10 additions & 3 deletions controllers/dataprotection/restore_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -54,16 +54,20 @@ import (
// RestoreReconciler reconciles a Restore object
type RestoreReconciler struct {
client.Client
Scheme *runtime.Scheme
Recorder record.EventRecorder
APIReader client.Reader
Scheme *runtime.Scheme
Recorder record.EventRecorder
}

// +kubebuilder:rbac:groups=dataprotection.kubeblocks.io,resources=restores,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=dataprotection.kubeblocks.io,resources=restores/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=dataprotection.kubeblocks.io,resources=restores/finalizers,verbs=update
// +kubebuilder:rbac:groups=batch,resources=jobs,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=core,resources=pods,verbs=get;list;watch
// +kubebuilder:rbac:groups=core,resources=serviceaccounts,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=rbac.authorization.k8s.io,resources=rolebindings,verbs=get;list;watch;create;update;patch;delete
// +kubebuilder:rbac:groups=apps.kubeblocks.io,resources=clusters;components,verbs=get;list;watch
// +kubebuilder:rbac:groups=workloads.kubeblocks.io,resources=instances;instancesets,verbs=get;list;watch

func (r *RestoreReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
reqCtx := intctrlutil.RequestCtx{
Expand Down Expand Up @@ -102,6 +106,9 @@ func (r *RestoreReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ct

// SetupWithManager sets up the controller with the Manager.
func (r *RestoreReconciler) SetupWithManager(mgr ctrl.Manager) error {
if r.APIReader == nil {
r.APIReader = mgr.GetAPIReader()
}
return intctrlutil.NewControllerManagedBy(mgr).
For(&dpv1alpha1.Restore{}).
Owns(&batchv1.Job{}).
Expand Down Expand Up @@ -498,7 +505,7 @@ func (r *RestoreReconciler) handleBackupActionSet(reqCtx intctrlutil.RequestCtx,
case dpv1alpha1.PostReady:
if len(jobs) == 0 {
// 2. build jobs for postReady action
jobs, err = restoreMgr.BuildPostReadyActionJobs(reqCtx, r.Client, backupSet, target, step)
jobs, err = restoreMgr.BuildPostReadyActionJobs(reqCtx, r.Client, r.APIReader, backupSet, target, step)
}
}
if err != nil {
Expand Down
94 changes: 63 additions & 31 deletions controllers/dataprotection/volumepopulator_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ type pvcRestoreDecision struct {

// +kubebuilder:rbac:groups=core,resources=persistentvolumeclaims/status,verbs=get;update;patch
// +kubebuilder:rbac:groups=core,resources=persistentvolumeclaims/finalizers,verbs=update
// +kubebuilder:rbac:groups=apps.kubeblocks.io,resources=clusters,verbs=get;list;watch
// +kubebuilder:rbac:groups=apps.kubeblocks.io,resources=components,verbs=get;list;watch
// +kubebuilder:rbac:groups=apps.kubeblocks.io,resources=componentdefinitions,verbs=get;list;watch

Expand Down Expand Up @@ -1172,8 +1173,37 @@ func (r *VolumePopulatorReconciler) ensurePostReadyRestoreCompleted(reqCtx intct
}
return false, err
}
if comp.Status.Phase != appsv1.RunningComponentPhase || componentPostProvisionRunning(comp) {
if err = r.updatePVCConditionsIfPopulateNotReleased(reqCtx, pvc, "Waiting for component to finish post-provision"); err != nil {
existing := &dpv1alpha1.Restore{}
existingKey := types.NamespacedName{
Namespace: pvc.Namespace,
Name: postReadyRestoreName(comp.UID),
}
if err = r.Client.Get(reqCtx.Ctx, existingKey, existing); err == nil {
desired, buildErr := r.buildPostReadyRestore(reqCtx, pvc, restoreMgr, comp, componentName, backupTarget)
if buildErr != nil {
return false, buildErr
}
if validateErr := validatePostReadyRestore(existing, desired, comp); validateErr != nil {
return false, validateErr
}
switch existing.Status.Phase {
case dpv1alpha1.RestorePhaseCompleted:
return true, nil
case dpv1alpha1.RestorePhaseFailed:
return false, intctrlutil.NewFatalError(fmt.Sprintf("postReady restore %s/%s failed", existing.Namespace, existing.Name))
default:
if err = r.updatePVCConditionsIfPopulateNotReleased(reqCtx, pvc, "Waiting for postReady restore to complete"); err != nil {
return false, err
}
return false, nil
}
} else if !apierrors.IsNotFound(err) {
return false, err
}
if comp.Generation != comp.Status.ObservedGeneration ||
Comment thread
weicao marked this conversation as resolved.
comp.Status.Phase != appsv1.RunningComponentPhase || componentPostProvisionRunning(comp) {
if err = r.updatePVCConditionsIfPopulateNotReleased(reqCtx, pvc,
"Waiting for component to observe the current generation and finish post-provision"); err != nil {
return false, err
}
return false, nil
Expand All @@ -1185,7 +1215,11 @@ func (r *VolumePopulatorReconciler) ensurePostReadyRestoreCompleted(reqCtx intct
if parameters[dptypes.DeferPostReadyUntilClusterRunningParameterKey] == "true" {
cluster := &appsv1.Cluster{}
if err = r.Client.Get(reqCtx.Ctx, types.NamespacedName{Namespace: pvc.Namespace, Name: clusterName}, cluster); err != nil {
return false, err
if apierrors.IsForbidden(err) || apierrors.IsUnauthorized(err) {
return false, intctrlutil.NewFatalError(err.Error())
}
return false, intctrlutil.NewRequeueError(reconcileInterval,
fmt.Sprintf("waiting for target cluster %s/%s: %v", pvc.Namespace, clusterName, err))
}
if cluster.Status.Phase != appsv1.RunningClusterPhase {
if err = r.updatePVCConditionsIfPopulateNotReleased(reqCtx, pvc, "Waiting for cluster to run before postReady restore"); err != nil {
Expand All @@ -1198,33 +1232,13 @@ func (r *VolumePopulatorReconciler) ensurePostReadyRestoreCompleted(reqCtx intct
if err != nil {
return false, err
}
existing := &dpv1alpha1.Restore{}
if err = r.Client.Get(reqCtx.Ctx, client.ObjectKeyFromObject(postReadyRestore), existing); err != nil {
if !apierrors.IsNotFound(err) {
return false, err
}
if err = r.Client.Create(reqCtx.Ctx, postReadyRestore); err != nil && !apierrors.IsAlreadyExists(err) {
return false, err
}
if err = r.updatePVCConditionsIfPopulateNotReleased(reqCtx, pvc, "Waiting for postReady restore to complete"); err != nil {
return false, err
}
return false, nil
}
if err = validatePostReadyRestore(existing, postReadyRestore, comp); err != nil {
if err = r.Client.Create(reqCtx.Ctx, postReadyRestore); err != nil && !apierrors.IsAlreadyExists(err) {
return false, err
}
switch existing.Status.Phase {
case dpv1alpha1.RestorePhaseCompleted:
return true, nil
case dpv1alpha1.RestorePhaseFailed:
return false, intctrlutil.NewFatalError(fmt.Sprintf("postReady restore %s/%s failed", existing.Namespace, existing.Name))
default:
if err = r.updatePVCConditionsIfPopulateNotReleased(reqCtx, pvc, "Waiting for postReady restore to complete"); err != nil {
return false, err
}
return false, nil
if err = r.updatePVCConditionsIfPopulateNotReleased(reqCtx, pvc, "Waiting for postReady restore to complete"); err != nil {
return false, err
}
return false, nil
}

func (r *VolumePopulatorReconciler) updatePVCConditionsIfPopulateNotReleased(reqCtx intctrlutil.RequestCtx,
Expand Down Expand Up @@ -1322,7 +1336,7 @@ func (r *VolumePopulatorReconciler) buildPostReadyRestore(reqCtx intctrlutil.Req
Namespace: backupNamespace,
},
RestoreTime: pvc.Annotations[constant.RestorePITRAnnotationKey],
Env: restoreMgr.Restore.Spec.Env,
Env: stripPostReadyTargetEnv(restoreMgr.Restore.Spec.Env),
Parameters: restoreParametersToPairs(restoreActionParameters(parameters)),
ReadyConfig: readyConfig,
},
Expand All @@ -1336,12 +1350,31 @@ func (r *VolumePopulatorReconciler) buildPostReadyRestore(reqCtx intctrlutil.Req
return restore, nil
}

func stripPostReadyTargetEnv(source []corev1.EnvVar) []corev1.EnvVar {
var env []corev1.EnvVar
for i := range source {
switch source[i].Name {
case dptypes.DPTargetClusterTopology,
dptypes.DPTargetComponentServiceVersion,
dptypes.DPTargetComponentServiceVersionSelector:
continue
default:
env = append(env, source[i])
}
}
return env
}

func validatePostReadyRestore(existing, desired *dpv1alpha1.Restore, comp *appsv1.Component) error {
if !hasOwnerReference(existing.OwnerReferences, comp.UID) {
return intctrlutil.NewFatalError(fmt.Sprintf("postReady restore %s/%s is not owned by component %s/%s",
existing.Namespace, existing.Name, comp.Namespace, comp.Name))
}
if !reflect.DeepEqual(existing.Spec, desired.Spec) {
existingSpec := existing.Spec.DeepCopy()
existingSpec.Env = stripPostReadyTargetEnv(existingSpec.Env)
desiredSpec := desired.Spec.DeepCopy()
desiredSpec.Env = stripPostReadyTargetEnv(desiredSpec.Env)
if !reflect.DeepEqual(existingSpec, desiredSpec) {
return intctrlutil.NewFatalError(fmt.Sprintf("postReady restore %s/%s spec does not match current restore intent",
existing.Namespace, existing.Name))
}
Expand Down Expand Up @@ -1610,8 +1643,7 @@ func postReadyRestoreLabels(pvc *corev1.PersistentVolumeClaim, comp *appsv1.Comp
}

func postReadyRestoreName(componentUID types.UID) string {
// Backup dataSource restore is an initial, single-attempt restore for a Component UID.
return constant.ShortenKubeName(fmt.Sprintf("restore-%s-post-ready", componentUID), constant.KubeNameMaxLength)
return dprestore.PostReadyRestoreName(componentUID)
}

func (r *VolumePopulatorReconciler) Cleanup(reqCtx intctrlutil.RequestCtx, pvc *corev1.PersistentVolumeClaim) error {
Expand Down
Loading
Loading