diff --git a/changelogs/unreleased/5760-Lyndon-Li b/changelogs/unreleased/5760-Lyndon-Li new file mode 100644 index 000000000..a71ebed28 --- /dev/null +++ b/changelogs/unreleased/5760-Lyndon-Li @@ -0,0 +1 @@ +Fix issue 5043, after the restore pod is scheduled, check if the node-agent pod is running in the same node. \ No newline at end of file diff --git a/pkg/cmd/server/server.go b/pkg/cmd/server/server.go index 65414b961..d4d73a5ce 100644 --- a/pkg/cmd/server/server.go +++ b/pkg/cmd/server/server.go @@ -680,7 +680,7 @@ func (s *server) runControllers(defaultVolumeSnapshotLocations map[string]string s.config.restoreResourcePriorities, s.kubeClient.CoreV1().Namespaces(), podvolume.NewRestorerFactory(s.repoLocker, s.repoEnsurer, s.veleroClient, s.kubeClient.CoreV1(), - s.sharedInformerFactory.Velero().V1().BackupRepositories().Informer().HasSynced, s.logger), + s.kubeClient.CoreV1(), s.kubeClient, s.sharedInformerFactory.Velero().V1().BackupRepositories().Informer().HasSynced, s.logger), s.config.podVolumeOperationTimeout, s.config.resourceTerminatingTimeout, s.logger, diff --git a/pkg/podvolume/restorer.go b/pkg/podvolume/restorer.go index 1360cdc65..c78a85d69 100644 --- a/pkg/podvolume/restorer.go +++ b/pkg/podvolume/restorer.go @@ -19,19 +19,24 @@ package podvolume import ( "context" "sync" + "time" "github.com/pkg/errors" "github.com/sirupsen/logrus" corev1api "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/client-go/kubernetes" corev1client "k8s.io/client-go/kubernetes/typed/core/v1" "k8s.io/client-go/tools/cache" velerov1api "github.com/vmware-tanzu/velero/pkg/apis/velero/v1" clientset "github.com/vmware-tanzu/velero/pkg/generated/clientset/versioned" "github.com/vmware-tanzu/velero/pkg/label" + "github.com/vmware-tanzu/velero/pkg/nodeagent" "github.com/vmware-tanzu/velero/pkg/repository" "github.com/vmware-tanzu/velero/pkg/util/boolptr" + "github.com/vmware-tanzu/velero/pkg/util/kube" ) type RestoreData struct { @@ -53,10 +58,13 @@ type restorer struct { repoEnsurer *repository.RepositoryEnsurer veleroClient clientset.Interface pvcClient corev1client.PersistentVolumeClaimsGetter + podClient corev1client.PodsGetter + kubeClient kubernetes.Interface - resultsLock sync.Mutex - results map[string]chan *velerov1api.PodVolumeRestore - log logrus.FieldLogger + resultsLock sync.Mutex + results map[string]chan *velerov1api.PodVolumeRestore + nodeAgentCheck chan error + log logrus.FieldLogger } func newRestorer( @@ -66,6 +74,8 @@ func newRestorer( podVolumeRestoreInformer cache.SharedIndexInformer, veleroClient clientset.Interface, pvcClient corev1client.PersistentVolumeClaimsGetter, + podClient corev1client.PodsGetter, + kubeClient kubernetes.Interface, log logrus.FieldLogger, ) *restorer { r := &restorer{ @@ -74,6 +84,8 @@ func newRestorer( repoEnsurer: repoEnsurer, veleroClient: veleroClient, pvcClient: pvcClient, + podClient: podClient, + kubeClient: kubeClient, results: make(map[string]chan *velerov1api.PodVolumeRestore), log: log, @@ -108,6 +120,10 @@ func (r *restorer) RestorePodVolumes(data RestoreData) []error { return nil } + if err := nodeagent.IsRunning(r.ctx, r.kubeClient, data.Restore.Namespace); err != nil { + return []error{errors.Wrapf(err, "error to check node agent status")} + } + repositoryType, err := getVolumesRepositoryType(volumesToRestore) if err != nil { return []error{err} @@ -129,6 +145,8 @@ func (r *restorer) RestorePodVolumes(data RestoreData) []error { r.results[resultsKey(data.Pod.Namespace, data.Pod.Name)] = resultsChan r.resultsLock.Unlock() + r.nodeAgentCheck = make(chan error) + var ( errs []error numRestores int @@ -161,6 +179,40 @@ func (r *restorer) RestorePodVolumes(data RestoreData) []error { numRestores++ } + checkCtx, checkCancel := context.WithCancel(context.Background()) + go func() { + nodeName := "" + + checkFunc := func(ctx context.Context) (bool, error) { + newObj, err := r.kubeClient.CoreV1().Pods(data.Pod.Namespace).Get(ctx, data.Pod.Name, metav1.GetOptions{}) + if err != nil { + return false, err + } + + nodeName = newObj.Spec.NodeName + + err = kube.IsPodScheduled(newObj) + if err != nil { + return false, nil + } else { + return true, nil + } + } + + err := wait.PollWithContext(checkCtx, time.Millisecond*500, time.Minute*10, checkFunc) + if err == wait.ErrWaitTimeout { + r.log.WithError(err).Error("Restoring pod is not scheduled until timeout or cancel, disengage") + } else if err != nil { + r.log.WithError(err).Error("Failed to check node-agent pod status, disengage") + } else { + err = nodeagent.IsRunningInNode(checkCtx, data.Restore.Namespace, nodeName, r.podClient) + if err != nil { + r.log.WithField("node", nodeName).WithError(err).Error("node-agent pod is not running in node, abort the restore") + r.nodeAgentCheck <- errors.Wrapf(err, "node-agent pod is not running in node %s", nodeName) + } + } + }() + ForEachVolume: for i := 0; i < numRestores; i++ { select { @@ -171,9 +223,17 @@ ForEachVolume: if res.Status.Phase == velerov1api.PodVolumeRestorePhaseFailed { errs = append(errs, errors.Errorf("pod volume restore failed: %s", res.Status.Message)) } + case err := <-r.nodeAgentCheck: + errs = append(errs, err) + break ForEachVolume } } + // This is to prevent the case that resultsChan is signaled before nodeAgentCheck though this is unlikely possible. + // One possible case is that the CR is edited and set to an ending state manually, either completed or failed. + // In this case, we must notify the check routine to stop. + checkCancel() + r.resultsLock.Lock() delete(r.results, resultsKey(data.Pod.Namespace, data.Pod.Name)) r.resultsLock.Unlock() diff --git a/pkg/podvolume/restorer_factory.go b/pkg/podvolume/restorer_factory.go index 3432f4be3..e4f2ce9c1 100644 --- a/pkg/podvolume/restorer_factory.go +++ b/pkg/podvolume/restorer_factory.go @@ -23,6 +23,7 @@ import ( "github.com/pkg/errors" "github.com/sirupsen/logrus" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/client-go/kubernetes" corev1client "k8s.io/client-go/kubernetes/typed/core/v1" "k8s.io/client-go/tools/cache" @@ -42,6 +43,8 @@ func NewRestorerFactory(repoLocker *repository.RepoLocker, repoEnsurer *repository.RepositoryEnsurer, veleroClient clientset.Interface, pvcClient corev1client.PersistentVolumeClaimsGetter, + podClient corev1client.PodsGetter, + kubeClient kubernetes.Interface, repoInformerSynced cache.InformerSynced, log logrus.FieldLogger) RestorerFactory { return &restorerFactory{ @@ -49,6 +52,8 @@ func NewRestorerFactory(repoLocker *repository.RepoLocker, repoEnsurer: repoEnsurer, veleroClient: veleroClient, pvcClient: pvcClient, + podClient: podClient, + kubeClient: kubeClient, repoInformerSynced: repoInformerSynced, log: log, } @@ -59,6 +64,8 @@ type restorerFactory struct { repoEnsurer *repository.RepositoryEnsurer veleroClient clientset.Interface pvcClient corev1client.PersistentVolumeClaimsGetter + podClient corev1client.PodsGetter + kubeClient kubernetes.Interface repoInformerSynced cache.InformerSynced log logrus.FieldLogger } @@ -74,7 +81,7 @@ func (rf *restorerFactory) NewRestorer(ctx context.Context, restore *velerov1api }, ) - r := newRestorer(ctx, rf.repoLocker, rf.repoEnsurer, informer, rf.veleroClient, rf.pvcClient, rf.log) + r := newRestorer(ctx, rf.repoLocker, rf.repoEnsurer, informer, rf.veleroClient, rf.pvcClient, rf.podClient, rf.kubeClient, rf.log) go informer.Run(ctx.Done()) if !cache.WaitForCacheSync(ctx.Done(), informer.HasSynced, rf.repoInformerSynced) { diff --git a/pkg/util/kube/pod.go b/pkg/util/kube/pod.go index 4e589c4f0..5d0b0f7f0 100644 --- a/pkg/util/kube/pod.go +++ b/pkg/util/kube/pod.go @@ -23,12 +23,38 @@ import ( // IsPodRunning does a well-rounded check to make sure the specified pod is running stably. // If not, return the error found func IsPodRunning(pod *corev1api.Pod) error { + return isPodScheduledInStatus(pod, func(pod *corev1api.Pod) error { + if pod.Status.Phase != corev1api.PodRunning { + return errors.New("pod is not running") + } + + return nil + }) +} + +// IsPodScheduled does a well-rounded check to make sure the specified pod has been scheduled into a node and in a stable status. +// If not, return the error found +func IsPodScheduled(pod *corev1api.Pod) error { + return isPodScheduledInStatus(pod, func(pod *corev1api.Pod) error { + if pod.Status.Phase != corev1api.PodRunning && pod.Status.Phase != corev1api.PodPending { + return errors.New("pod is running or pending") + } + + return nil + }) +} + +func isPodScheduledInStatus(pod *corev1api.Pod, statusCheckFunc func(*corev1api.Pod) error) error { + if pod == nil { + return errors.New("invalid input pod") + } + if pod.Spec.NodeName == "" { return errors.Errorf("pod is not scheduled, name=%s, namespace=%s, phase=%s", pod.Name, pod.Namespace, pod.Status.Phase) } - if pod.Status.Phase != corev1api.PodRunning { - return errors.Errorf("pod is not running, name=%s, namespace=%s, phase=%s", pod.Name, pod.Namespace, pod.Status.Phase) + if err := statusCheckFunc(pod); err != nil { + return errors.Wrapf(err, "pod is not in the expected status, name=%s, namespace=%s, phase=%s", pod.Name, pod.Namespace, pod.Status.Phase) } if pod.DeletionTimestamp != nil {