Skip to content
Draft
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
24 changes: 13 additions & 11 deletions pkg/dataprotection/restore/builder.go
Original file line number Diff line number Diff line change
Expand Up @@ -253,9 +253,9 @@ func (r *restoreJobBuilder) addCommonEnv(sourceTargetPodName string) *restoreJob

func (r *restoreJobBuilder) addTargetPodAndCredentialEnv(pod *corev1.Pod,
connectionCredential *dpv1alpha1.ConnectionCredential,
target *dpv1alpha1.BackupTarget) *restoreJobBuilder {
target *dpv1alpha1.BackupTarget) (*restoreJobBuilder, error) {
if pod == nil {
return r
return r, nil
}
var env []corev1.EnvVar
// Note: now only add the first container envs.
Expand All @@ -266,19 +266,19 @@ func (r *restoreJobBuilder) addTargetPodAndCredentialEnv(pod *corev1.Pod,
addDBHostEnv := func() {
env = append(env, corev1.EnvVar{Name: dptypes.DPDBHost, Value: intctrlutil.BuildPodHostDNS(pod)})
}
addDBPortEnv := func() {
addDBPortEnv := func() error {
portEnv, err := utils.GetDPDBPortEnv(pod, target.ContainerPort)
if err != nil {
// fallback to use the first port of the pod
portEnv, _ = utils.GetDPDBPortEnv(pod, nil)
}
if portEnv != nil {
env = append(env, *portEnv)
return err
}
env = append(env, *portEnv)
return nil
}
if connectionCredential == nil {
addDBHostEnv()
addDBPortEnv()
if err := addDBPortEnv(); err != nil {
return r, err
}
} else {
appendEnvFromSecret := func(envName, keyName string) {
if keyName == "" {
Expand All @@ -298,7 +298,9 @@ func (r *restoreJobBuilder) addTargetPodAndCredentialEnv(pod *corev1.Pod,
if connectionCredential.PortKey != "" {
appendEnvFromSecret(dptypes.DPDBPort, connectionCredential.PortKey)
} else {
addDBPortEnv()
if err := addDBPortEnv(); err != nil {
return r, err
}
}
if connectionCredential.HostKey != "" {
appendEnvFromSecret(dptypes.DPDBHost, connectionCredential.HostKey)
Expand All @@ -307,7 +309,7 @@ func (r *restoreJobBuilder) addTargetPodAndCredentialEnv(pod *corev1.Pod,
}
}
r.env = utils.MergeEnv(r.env, env)
return r
return r, nil
}

// builderRestoreJobName builds restore job name.
Expand Down
19 changes: 13 additions & 6 deletions pkg/dataprotection/restore/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -720,7 +720,7 @@ func (r *RestoreManager) BuildPostReadyActionJobs(reqCtx intctrlutil.RequestCtx,
return nil, err
}
sort.Sort(intctrlutil.ByPodName(targetPodList.Items))
buildJob := func(targetPod *corev1.Pod, sourceTargetPodName string, index int) *batchv1.Job {
buildJob := func(targetPod *corev1.Pod, sourceTargetPodName string, index int) (*batchv1.Job, error) {
if boolptr.IsSetToTrue(actionSpec.Job.RunOnTargetPodNode) {
jobBuilder.resetSpecificVolumesAndMounts()
jobBuilder.setNodeNameToNodeSelector(targetPod.Spec.NodeName)
Expand All @@ -734,15 +734,18 @@ func (r *RestoreManager) BuildPostReadyActionJobs(reqCtx intctrlutil.RequestCtx,
}
}
}
return jobBuilder.setImage(actionSpec.Job.Image).
jobBuilder.setImage(actionSpec.Job.Image).
setJobName(buildJobName(index)).
addCommonEnv(sourceTargetPodName).
attachBackupRepo().
setCommand(actionSpec.Job.Command).
setToleration(targetPod.Spec.Tolerations).
addTargetPodAndCredentialEnv(targetPod, readyConfig.ConnectionCredential, &target.BackupTarget).
setToleration(targetPod.Spec.Tolerations)
if _, err := jobBuilder.addTargetPodAndCredentialEnv(targetPod, readyConfig.ConnectionCredential, &target.BackupTarget); err != nil {
return nil, err
}
return jobBuilder.
setServiceAccount(r.WorkerServiceAccount).
build()
build(), nil
}

if podSelector.Strategy == dpv1alpha1.PodSelectionStrategyAny {
Expand All @@ -762,7 +765,11 @@ func (r *RestoreManager) BuildPostReadyActionJobs(reqCtx intctrlutil.RequestCtx,
// no need to recover the volume when the pod selection policy is 'All' and sourceTargetPodName is not found.
continue
}
jobs = append(jobs, buildJob(&targetPodList.Items[i], sourceTargetPodName, i))
job, err := buildJob(&targetPodList.Items[i], sourceTargetPodName, i)
if err != nil {
return nil, err
}
jobs = append(jobs, job)
}
return jobs, nil
}
Expand Down
48 changes: 48 additions & 0 deletions pkg/dataprotection/restore/utils_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,54 @@ func TestRestoreBuilderCommonVolumesAndSourceMountPath(t *testing.T) {
assert.Len(t, builder.commonVolumes, 1)
}

func TestRestoreTargetPortLookupDoesNotFallbackToFirstPodPort(t *testing.T) {
pod := &corev1.Pod{
Spec: corev1.PodSpec{Containers: []corev1.Container{
{
Name: "fe",
Ports: []corev1.ContainerPort{
{Name: "http-port", ContainerPort: 8030},
{Name: "query-port", ContainerPort: 9030},
},
},
}},
}
credential := &dpv1alpha1.ConnectionCredential{SecretName: "connection"}

builder := &restoreJobBuilder{}
_, err := builder.addTargetPodAndCredentialEnv(pod, credential, &dpv1alpha1.BackupTarget{
ContainerPort: &dpv1alpha1.ContainerPort{ContainerName: "fe", PortName: "query-port"},
})
assert.NoError(t, err)
assert.Contains(t, builder.env, corev1.EnvVar{Name: dptypes.DPDBPort, Value: "9030"})

builder = &restoreJobBuilder{}
_, err = builder.addTargetPodAndCredentialEnv(pod, credential, &dpv1alpha1.BackupTarget{})
assert.NoError(t, err)
assert.Contains(t, builder.env, corev1.EnvVar{Name: dptypes.DPDBPort, Value: "8030"})

builder = &restoreJobBuilder{}
_, err = builder.addTargetPodAndCredentialEnv(pod, credential, &dpv1alpha1.BackupTarget{
ContainerPort: &dpv1alpha1.ContainerPort{ContainerName: "fe", PortName: "missing"},
})
assert.ErrorContains(t, err, "specified containerPort")
assert.NotContains(t, builder.env, corev1.EnvVar{Name: dptypes.DPDBPort, Value: "8030"})

builder = &restoreJobBuilder{}
credential.PortKey = "port"
_, err = builder.addTargetPodAndCredentialEnv(pod, credential, &dpv1alpha1.BackupTarget{
ContainerPort: &dpv1alpha1.ContainerPort{ContainerName: "fe", PortName: "missing"},
})
assert.NoError(t, err)
assert.Contains(t, builder.env, corev1.EnvVar{
Name: dptypes.DPDBPort,
ValueFrom: &corev1.EnvVarSource{SecretKeyRef: &corev1.SecretKeySelector{
LocalObjectReference: corev1.LocalObjectReference{Name: "connection"},
Key: "port",
}},
})
}

func TestRestoreConditionAndStatusHelpers(t *testing.T) {
restore := &dpv1alpha1.Restore{}
SetRestoreCheckBackupRepoCondition(restore, ReasonCheckBackupRepoSuccessfully, "ok")
Expand Down
Loading