The backup and restore data mover node selection.

Update Makefile to let the `make serve-docs` work again.

Signed-off-by: Xun Jiang <xun.jiang@broadcom.com>
This commit is contained in:
Xun Jiang
2025-07-01 16:26:12 +08:00
parent bd609db395
commit f2133c7d22
18 changed files with 851 additions and 211 deletions
+18 -3
View File
@@ -64,6 +64,7 @@ type DataDownloadReconciler struct {
restoreExposer exposer.GenericRestoreExposer
nodeName string
dataPathMgr *datapath.Manager
loadAffinity []*kube.LoadAffinity
restorePVCConfig nodeagent.RestorePVC
podResources corev1api.ResourceRequirements
preparingTimeout time.Duration
@@ -71,9 +72,19 @@ type DataDownloadReconciler struct {
cancelledDataDownload map[string]time.Time
}
func NewDataDownloadReconciler(client client.Client, mgr manager.Manager, kubeClient kubernetes.Interface, dataPathMgr *datapath.Manager,
restorePVCConfig nodeagent.RestorePVC, podResources corev1api.ResourceRequirements, nodeName string, preparingTimeout time.Duration,
logger logrus.FieldLogger, metrics *metrics.ServerMetrics) *DataDownloadReconciler {
func NewDataDownloadReconciler(
client client.Client,
mgr manager.Manager,
kubeClient kubernetes.Interface,
dataPathMgr *datapath.Manager,
loadAffinity []*kube.LoadAffinity,
restorePVCConfig nodeagent.RestorePVC,
podResources corev1api.ResourceRequirements,
nodeName string,
preparingTimeout time.Duration,
logger logrus.FieldLogger,
metrics *metrics.ServerMetrics,
) *DataDownloadReconciler {
return &DataDownloadReconciler{
client: client,
kubeClient: kubeClient,
@@ -84,6 +95,7 @@ func NewDataDownloadReconciler(client client.Client, mgr manager.Manager, kubeCl
restoreExposer: exposer.NewGenericRestoreExposer(kubeClient, logger),
restorePVCConfig: restorePVCConfig,
dataPathMgr: dataPathMgr,
loadAffinity: loadAffinity,
podResources: podResources,
preparingTimeout: preparingTimeout,
metrics: metrics,
@@ -828,6 +840,8 @@ func (r *DataDownloadReconciler) setupExposeParam(dd *velerov2alpha1api.DataDown
}
}
affinity := kube.GetLoadAffinityByStorageClass(r.loadAffinity, dd.Spec.BackupStorageLocation, log)
return exposer.GenericRestoreExposeParam{
TargetPVCName: dd.Spec.TargetVolume.PVC,
TargetNamespace: dd.Spec.TargetVolume.Namespace,
@@ -838,6 +852,7 @@ func (r *DataDownloadReconciler) setupExposeParam(dd *velerov2alpha1api.DataDown
ExposeTimeout: r.preparingTimeout,
NodeOS: nodeOS,
RestorePVCConfig: r.restorePVCConfig,
LoadAffinity: affinity,
}, nil
}
+45 -43
View File
@@ -40,8 +40,6 @@ import (
"sigs.k8s.io/controller-runtime/pkg/manager"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
"sigs.k8s.io/controller-runtime/pkg/client/fake"
velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1"
velerov2alpha1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v2alpha1"
"github.com/vmware-tanzu/velero/pkg/builder"
@@ -70,7 +68,9 @@ func dataDownloadBuilder() *builder.DataDownloadBuilder {
})
}
func initDataDownloadReconciler(objects []runtime.Object, needError ...bool) (*DataDownloadReconciler, error) {
func initDataDownloadReconciler(t *testing.T, objects []any, needError ...bool) (*DataDownloadReconciler, error) {
t.Helper()
var errs = make([]error, 6)
for k, isError := range needError {
if k == 0 && isError {
@@ -87,28 +87,24 @@ func initDataDownloadReconciler(objects []runtime.Object, needError ...bool) (*D
errs[5] = fmt.Errorf("List error")
}
}
return initDataDownloadReconcilerWithError(objects, errs...)
return initDataDownloadReconcilerWithError(t, objects, errs...)
}
func initDataDownloadReconcilerWithError(objects []runtime.Object, needError ...error) (*DataDownloadReconciler, error) {
scheme := runtime.NewScheme()
err := velerov1api.AddToScheme(scheme)
if err != nil {
return nil, err
}
err = velerov2alpha1api.AddToScheme(scheme)
if err != nil {
return nil, err
}
err = corev1api.AddToScheme(scheme)
if err != nil {
return nil, err
func initDataDownloadReconcilerWithError(t *testing.T, objects []any, needError ...error) (*DataDownloadReconciler, error) {
t.Helper()
runtimeObjects := make([]runtime.Object, 0)
for _, obj := range objects {
runtimeObjects = append(runtimeObjects, obj.(runtime.Object))
}
fakeClient := &FakeClient{
Client: fake.NewClientBuilder().WithScheme(scheme).Build(),
fakeClient := FakeClient{
Client: velerotest.NewFakeControllerRuntimeClient(t, runtimeObjects...),
}
fakeKubeClient := clientgofake.NewSimpleClientset(runtimeObjects...)
for k := range needError {
if k == 0 {
fakeClient.getError = needError[0]
@@ -125,26 +121,32 @@ func initDataDownloadReconcilerWithError(objects []runtime.Object, needError ...
}
}
var fakeKubeClient *clientgofake.Clientset
if len(objects) != 0 {
fakeKubeClient = clientgofake.NewSimpleClientset(objects...)
} else {
fakeKubeClient = clientgofake.NewSimpleClientset()
}
fakeFS := velerotest.NewFakeFileSystem()
pathGlob := fmt.Sprintf("/host_pods/%s/volumes/*/%s", "test-uid", "test-pvc")
_, err = fakeFS.Create(pathGlob)
_, err := fakeFS.Create(pathGlob)
if err != nil {
return nil, err
}
dataPathMgr := datapath.NewManager(1)
return NewDataDownloadReconciler(fakeClient, nil, fakeKubeClient, dataPathMgr, nodeagent.RestorePVC{}, corev1api.ResourceRequirements{}, "test-node", time.Minute*5, velerotest.NewLogger(), metrics.NewServerMetrics()), nil
return NewDataDownloadReconciler(
&fakeClient,
nil,
fakeKubeClient,
dataPathMgr,
nil,
nodeagent.RestorePVC{},
corev1api.ResourceRequirements{},
"test-node",
time.Minute*5,
velerotest.NewLogger(),
metrics.NewServerMetrics()), nil
}
func TestDataDownloadReconcile(t *testing.T) {
sc := builder.ForStorageClass("sc").Result()
daemonSet := &appsv1api.DaemonSet{
ObjectMeta: metav1.ObjectMeta{
Namespace: "velero",
@@ -330,7 +332,7 @@ func TestDataDownloadReconcile(t *testing.T) {
{
name: "dd succeeds for accepted",
dd: dataDownloadBuilder().Finalizers([]string{DataUploadDownloadFinalizer}).Result(),
targetPVC: builder.ForPersistentVolumeClaim("test-ns", "test-pvc").Result(),
targetPVC: builder.ForPersistentVolumeClaim("test-ns", "test-pvc").StorageClass("sc").Result(),
expected: dataDownloadBuilder().Finalizers([]string{DataUploadDownloadFinalizer}).Phase(velerov2alpha1api.DataDownloadPhaseAccepted).Result(),
},
{
@@ -457,13 +459,13 @@ func TestDataDownloadReconcile(t *testing.T) {
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
objs := []runtime.Object{daemonSet, node}
objects := []any{daemonSet, node, sc}
if test.targetPVC != nil {
objs = append(objs, test.targetPVC)
objects = append(objects, test.targetPVC)
}
r, err := initDataDownloadReconciler(objs, test.needErrs...)
r, err := initDataDownloadReconciler(t, objects, test.needErrs...)
require.NoError(t, err)
if !test.notCreateDD {
@@ -607,7 +609,7 @@ func TestOnDataDownloadFailed(t *testing.T) {
for _, getErr := range []bool{true, false} {
ctx := context.TODO()
needErrs := []bool{getErr, false, false, false}
r, err := initDataDownloadReconciler(nil, needErrs...)
r, err := initDataDownloadReconciler(t, nil, needErrs...)
require.NoError(t, err)
dd := dataDownloadBuilder().Result()
@@ -633,7 +635,7 @@ func TestOnDataDownloadCancelled(t *testing.T) {
for _, getErr := range []bool{true, false} {
ctx := context.TODO()
needErrs := []bool{getErr, false, false, false}
r, err := initDataDownloadReconciler(nil, needErrs...)
r, err := initDataDownloadReconciler(t, nil, needErrs...)
require.NoError(t, err)
dd := dataDownloadBuilder().Result()
@@ -675,7 +677,7 @@ func TestOnDataDownloadCompleted(t *testing.T) {
t.Run(test.name, func(t *testing.T) {
ctx := context.TODO()
needErrs := []bool{test.isGetErr, false, false, false}
r, err := initDataDownloadReconciler(nil, needErrs...)
r, err := initDataDownloadReconciler(t, nil, needErrs...)
r.restoreExposer = func() exposer.GenericRestoreExposer {
ep := exposermockes.NewGenericRestoreExposer(t)
if test.rebindVolumeErr {
@@ -740,7 +742,7 @@ func TestOnDataDownloadProgress(t *testing.T) {
t.Run(test.name, func(t *testing.T) {
ctx := context.TODO()
r, err := initDataDownloadReconciler(nil, test.needErrs...)
r, err := initDataDownloadReconciler(t, nil, test.needErrs...)
require.NoError(t, err)
defer func() {
r.client.Delete(ctx, test.dd, &kbclient.DeleteOptions{})
@@ -774,7 +776,7 @@ func TestOnDataDownloadProgress(t *testing.T) {
func TestFindDataDownloadForPod(t *testing.T) {
needErrs := []bool{false, false, false, false}
r, err := initDataDownloadReconciler(nil, needErrs...)
r, err := initDataDownloadReconciler(t, nil, needErrs...)
require.NoError(t, err)
tests := []struct {
name string
@@ -860,7 +862,7 @@ func TestAcceptDataDownload(t *testing.T) {
}
for _, test := range tests {
ctx := context.Background()
r, err := initDataDownloadReconcilerWithError(nil, test.needErrs...)
r, err := initDataDownloadReconcilerWithError(t, nil, test.needErrs...)
require.NoError(t, err)
err = r.client.Create(ctx, test.dd)
@@ -904,7 +906,7 @@ func TestOnDdPrepareTimeout(t *testing.T) {
}
for _, test := range tests {
ctx := context.Background()
r, err := initDataDownloadReconcilerWithError(nil, test.needErrs...)
r, err := initDataDownloadReconcilerWithError(t, nil, test.needErrs...)
require.NoError(t, err)
err = r.client.Create(ctx, test.dd)
@@ -949,7 +951,7 @@ func TestTryCancelDataDownload(t *testing.T) {
}
for _, test := range tests {
ctx := context.Background()
r, err := initDataDownloadReconcilerWithError(nil, test.needErrs...)
r, err := initDataDownloadReconcilerWithError(t, nil, test.needErrs...)
require.NoError(t, err)
err = r.client.Create(ctx, test.dd)
@@ -1007,7 +1009,7 @@ func TestUpdateDataDownloadWithRetry(t *testing.T) {
t.Run(tc.Name, func(t *testing.T) {
ctx, cancelFunc := context.WithTimeout(context.TODO(), time.Second*5)
defer cancelFunc()
r, err := initDataDownloadReconciler(nil, tc.needErrs...)
r, err := initDataDownloadReconciler(t, nil, tc.needErrs...)
require.NoError(t, err)
err = r.client.Create(ctx, dataDownloadBuilder().Result())
require.NoError(t, err)
@@ -1124,7 +1126,7 @@ func TestAttemptDataDownloadResume(t *testing.T) {
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
ctx := context.TODO()
r, err := initDataDownloadReconciler(nil, test.needErrs...)
r, err := initDataDownloadReconciler(t, nil, test.needErrs...)
r.nodeName = "node-1"
require.NoError(t, err)
defer func() {
@@ -1242,7 +1244,7 @@ func TestResumeCancellableRestore(t *testing.T) {
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
ctx := context.TODO()
r, err := initDataDownloadReconciler(nil, false)
r, err := initDataDownloadReconciler(t, nil, false)
r.nodeName = "node-1"
require.NoError(t, err)
+5 -3
View File
@@ -74,7 +74,7 @@ type DataUploadReconciler struct {
logger logrus.FieldLogger
snapshotExposerList map[velerov2alpha1api.SnapshotType]exposer.SnapshotExposer
dataPathMgr *datapath.Manager
loadAffinity *kube.LoadAffinity
loadAffinity []*kube.LoadAffinity
backupPVCConfig map[string]nodeagent.BackupPVC
podResources corev1api.ResourceRequirements
preparingTimeout time.Duration
@@ -88,7 +88,7 @@ func NewDataUploadReconciler(
kubeClient kubernetes.Interface,
csiSnapshotClient snapshotter.SnapshotV1Interface,
dataPathMgr *datapath.Manager,
loadAffinity *kube.LoadAffinity,
loadAffinity []*kube.LoadAffinity,
backupPVCConfig map[string]nodeagent.BackupPVC,
podResources corev1api.ResourceRequirements,
clock clocks.WithTickerAndDelayedExecution,
@@ -917,6 +917,8 @@ func (r *DataUploadReconciler) setupExposeParam(du *velerov2alpha1api.DataUpload
}
}
affinity := kube.GetLoadAffinityByStorageClass(r.loadAffinity, du.Spec.CSISnapshot.SnapshotClass, log)
return &exposer.CSISnapshotExposeParam{
SnapshotName: du.Spec.CSISnapshot.VolumeSnapshot,
SourceNamespace: du.Spec.SourceNamespace,
@@ -927,7 +929,7 @@ func (r *DataUploadReconciler) setupExposeParam(du *velerov2alpha1api.DataUpload
OperationTimeout: du.Spec.OperationTimeout.Duration,
ExposeTimeout: r.preparingTimeout,
VolumeSize: pvc.Spec.Resources.Requests[corev1api.ResourceStorage],
Affinity: r.loadAffinity,
Affinity: affinity,
BackupPVCConfig: r.backupPVCConfig,
Resources: r.podResources,
NodeOS: nodeOS,