summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--cloud/pkg/cloudhub/channelq/channelq.go5
-rw-r--r--cloud/pkg/edgecontroller/controller/upstream.go2
-rw-r--r--edge/pkg/edged/edged.go124
3 files changed, 45 insertions, 86 deletions
diff --git a/cloud/pkg/cloudhub/channelq/channelq.go b/cloud/pkg/cloudhub/channelq/channelq.go
index cd9c5fa7f..3da0a0f9b 100644
--- a/cloud/pkg/cloudhub/channelq/channelq.go
+++ b/cloud/pkg/cloudhub/channelq/channelq.go
@@ -167,10 +167,7 @@ func getMsgKey(obj interface{}) (string, error) {
msg := obj.(*beehiveModel.Message)
if msg.GetGroup() == edgeconst.GroupResource {
- resourceType, _ := edgemessagelayer.GetResourceType(*msg)
- resourceNamespace, _ := edgemessagelayer.GetNamespace(*msg)
- resourceName, _ := edgemessagelayer.GetResourceName(*msg)
- return strings.Join([]string{resourceType, resourceNamespace, resourceName}, "/"), nil
+ return GetMessageUID(*msg)
}
return "", fmt.Errorf("failed to get message key")
diff --git a/cloud/pkg/edgecontroller/controller/upstream.go b/cloud/pkg/edgecontroller/controller/upstream.go
index 4fea9b162..c276cf458 100644
--- a/cloud/pkg/edgecontroller/controller/upstream.go
+++ b/cloud/pkg/edgecontroller/controller/upstream.go
@@ -297,7 +297,7 @@ func (uc *UpstreamController) updatePodStatus() {
case model.UpdateOperation:
for _, podStatus := range podStatuses {
getPod, err := uc.kubeClient.CoreV1().Pods(namespace).Get(context.Background(), podStatus.Name, metaV1.GetOptions{})
- if errors.IsNotFound(err) {
+ if (err == nil && getPod.UID != podStatus.UID) || errors.IsNotFound(err) {
klog.Warningf("message: %s, pod not found, namespace: %s, name: %s", msg.GetID(), namespace, podStatus.Name)
// Send request to delete this pod on edge side
diff --git a/edge/pkg/edged/edged.go b/edge/pkg/edged/edged.go
index 9df1e38ca..715bcde35 100644
--- a/edge/pkg/edged/edged.go
+++ b/edge/pkg/edged/edged.go
@@ -878,11 +878,7 @@ func (e *edged) syncLoopIteration(plegCh <-chan *pleg.PodLifecycleEvent, houseke
}
}
klog.Infof("Will restart pod [%s]", pod.Name)
- key := types.NamespacedName{
- Namespace: pod.Namespace,
- Name: pod.Name,
- }
- e.podAdditionQueue.Add(key.String())
+ e.podAdditionQueue.Add(pod.UID)
}
case update := <-e.readinessManager.Updates():
@@ -896,11 +892,7 @@ func (e *edged) syncLoopIteration(plegCh <-chan *pleg.PodLifecycleEvent, houseke
break
}
- key := types.NamespacedName{
- Namespace: pod.Namespace,
- Name: pod.Name,
- }
- e.podAdditionQueue.Add(key.String())
+ e.podAdditionQueue.Add(pod.UID)
case update := <-e.startupManager.Updates():
started := update.Result == proberesults.Success
@@ -913,11 +905,7 @@ func (e *edged) syncLoopIteration(plegCh <-chan *pleg.PodLifecycleEvent, houseke
break
}
- key := types.NamespacedName{
- Namespace: pod.Namespace,
- Name: pod.Name,
- }
- e.podAdditionQueue.Add(key.String())
+ e.podAdditionQueue.Add(pod.UID)
case plegEvent := <-plegCh:
if pod, ok := e.podManager.GetPodByUID(plegEvent.ID); ok {
@@ -942,11 +930,7 @@ func (e *edged) syncLoopIteration(plegCh <-chan *pleg.PodLifecycleEvent, houseke
}
}
klog.Infof("sync loop get event container died, restart pod [%s]", pod.Name)
- key := types.NamespacedName{
- Namespace: pod.Namespace,
- Name: pod.Name,
- }
- e.podAdditionQueue.Add(key.String())
+ e.podAdditionQueue.Add(pod.UID)
} else {
klog.Infof("sync loop get event [%s], ignore it now.", plegEvent.Type)
}
@@ -965,29 +949,13 @@ func (e *edged) syncLoopIteration(plegCh <-chan *pleg.PodLifecycleEvent, houseke
}
for _, pod := range podsToSync {
if !e.podIsTerminated(pod) {
- key := types.NamespacedName{
- Namespace: pod.Namespace,
- Name: pod.Name,
- }
- e.podAdditionQueue.Add(key.String())
+ e.podAdditionQueue.Add(pod.UID)
}
}
}
}
}
-// NewNamespacedNameFromString parses the provided string and returns a NamespacedName
-func NewNamespacedNameFromString(s string) types.NamespacedName {
- Separator := '/'
- nn := types.NamespacedName{}
- result := strings.Split(s, string(Separator))
- if len(result) == 2 {
- nn.Namespace = result[0]
- nn.Name = result[1]
- }
- return nn
-}
-
func (e *edged) podAddWorkerRun(consumers int) {
for i := 0; i < consumers; i++ {
klog.Infof("start pod addition queue work %d", i)
@@ -998,28 +966,28 @@ func (e *edged) podAddWorkerRun(consumers int) {
klog.Errorf("consumer: [%d], worker addition queue is shutting down!", i)
return
}
- namespacedName := NewNamespacedNameFromString(item.(string))
- podName := namespacedName.Name
- klog.Infof("worker [%d] get pod addition item [%s]", i, podName)
- backOffKey := fmt.Sprintf("pod_addition_worker_%s", podName)
+
+ podUID := item.(types.UID)
+ klog.Infof("worker [%d] get pod addition item [%s]", i, podUID)
+ backOffKey := fmt.Sprintf("pod_addition_worker_%s", podUID)
if e.podAdditionBackoff.IsInBackOffSinceUpdate(backOffKey, e.podAdditionBackoff.Clock.Now()) {
- klog.Errorf("consume pod addition backoff: Back-off consume pod [%s] addition error, backoff: [%v]", podName, e.podAdditionBackoff.Get(backOffKey))
+ klog.Errorf("consume pod addition backoff: Back-off consume pod [%s] addition error, backoff: [%v]", podUID, e.podAdditionBackoff.Get(backOffKey))
go func() {
- klog.Infof("worker [%d] backoff pod addition item [%s] failed, re-add to queue", i, podName)
+ klog.Infof("worker [%d] backoff pod addition item [%s] failed, re-add to queue", i, podUID)
time.Sleep(e.podAdditionBackoff.Get(backOffKey))
e.podAdditionQueue.Add(item)
}()
e.podAdditionQueue.Done(item)
continue
}
- err := e.consumePodAddition(&namespacedName)
+ err := e.consumePodAddition(podUID)
if err != nil {
if err == apis.ErrPodNotFound {
- klog.Errorf("worker [%d] handle pod addition item [%s] failed with not found error.", i, podName)
+ klog.Errorf("worker [%d] handle pod addition item [%s] failed with not found error.", i, podUID)
e.podAdditionBackoff.Reset(backOffKey)
} else {
go func() {
- klog.Errorf("worker [%d] handle pod addition item [%s] failed: %v, re-add to queue", i, podName, err)
+ klog.Errorf("worker [%d] handle pod addition item [%s] failed: %v, re-add to queue", i, podUID, err)
e.podAdditionBackoff.Next(backOffKey, e.podAdditionBackoff.Clock.Now())
time.Sleep(enqueueDuration)
e.podAdditionQueue.Add(item)
@@ -1043,18 +1011,19 @@ func (e *edged) podRemoveWorkerRun(consumers int) {
klog.Errorf("consumer: [%d], worker addition queue is shutting down!", i)
return
}
- namespacedName := NewNamespacedNameFromString(item.(string))
- podName := namespacedName.Name
- klog.Infof("consumer: [%d], worker get removed pod [%s]\n", i, podName)
- err := e.consumePodDeletion(&namespacedName)
+
+ podUID := item.(types.UID)
+
+ klog.Infof("consumer: [%d], worker get removed pod [%s]\n", i, podUID)
+ err := e.consumePodDeletion(podUID)
if err != nil {
if err == apis.ErrContainerNotFound {
- klog.Infof("pod [%s] is not exist, with container not found error", podName)
+ klog.Infof("pod [%s] is not exist, with container not found error", podUID)
} else if err == apis.ErrPodNotFound {
- klog.Infof("pod [%s] is not found", podName)
+ klog.Infof("pod [%s] is not found", podUID)
} else {
go func(item interface{}) {
- klog.Errorf("worker remove pod [%s] failed: %v", podName, err)
+ klog.Errorf("worker remove pod [%s] failed: %v", podUID, err)
time.Sleep(2 * time.Second)
e.podDeletionQueue.Add(item)
}(item)
@@ -1066,14 +1035,14 @@ func (e *edged) podRemoveWorkerRun(consumers int) {
}
}
-func (e *edged) consumePodAddition(namespacedName *types.NamespacedName) error {
- podName := namespacedName.Name
- klog.Infof("start to consume added pod [%s]", podName)
- pod, ok := e.podManager.GetPodByName(namespacedName.Namespace, podName)
+func (e *edged) consumePodAddition(podUID types.UID) error {
+ pod, ok := e.podManager.GetPodByUID(podUID)
if !ok || pod.DeletionTimestamp != nil {
return apis.ErrPodNotFound
}
+ klog.Infof("start to consume added pod [%s]", format.Pod(pod))
+
if err := e.makePodDataDirs(pod); err != nil {
return fmt.Errorf("unable to make pod data directories for pod %q: %v", format.Pod(pod), err)
}
@@ -1085,14 +1054,13 @@ func (e *edged) consumePodAddition(namespacedName *types.NamespacedName) error {
// Fetch the pull secrets for the pod
secrets := e.getImagePullSecretsForPod(pod)
- podUID := pod.GetUID()
t, ok := e.podLastSyncTime.Load(podUID)
if !ok {
t = time.Time{}
}
curPodStatus, err := e.podCache.GetNewerThan(podUID, t.(time.Time))
if err != nil {
- return fmt.Errorf("pod %s cache newer failed: %v", podName, err)
+ return fmt.Errorf("pod %s cache newer failed: %v", format.Pod(pod), err)
}
result := e.containerRuntime.SyncPod(pod, curPodStatus, secrets, e.podAdditionBackoff)
e.podLastSyncTime.Store(podUID, time.Now())
@@ -1101,8 +1069,9 @@ func (e *edged) consumePodAddition(namespacedName *types.NamespacedName) error {
for _, r := range result.SyncResults {
if r.Error != kubecontainer.ErrCrashLoopBackOff && r.Error != images.ErrImagePullBackOff {
// Do not record an event here, as we keep all event logging for sync pod failures
- // local to container runtime so we get better errors
- return fmt.Errorf("sync pod failed: %v", err)
+
+ // local to container runtime, so we get better errors
+ return fmt.Errorf("sync pod %s failed: %v", format.Pod(pod), err)
}
}
@@ -1110,21 +1079,21 @@ func (e *edged) consumePodAddition(namespacedName *types.NamespacedName) error {
}
e.workQueue.Enqueue(pod.UID, utilwait.Jitter(time.Minute, workerResyncIntervalJitterFactor))
- klog.Infof("consume added pod [%s] successfully\n", podName)
+ klog.Infof("consume added pod %s successfully\n", format.Pod(pod))
return nil
}
-func (e *edged) consumePodDeletion(namespacedName *types.NamespacedName) error {
- podName := namespacedName.Name
- klog.Infof("start to consume removed pod [%s]", podName)
- pod, ok := e.podManager.GetPodByName(namespacedName.Namespace, podName)
+func (e *edged) consumePodDeletion(podUID types.UID) error {
+ pod, ok := e.podManager.GetPodByUID(podUID)
if !ok {
return apis.ErrPodNotFound
}
+ klog.Infof("start to consume removed pod [%s]", format.Pod(pod))
+
podStatus, err := e.podCache.Get(pod.GetUID())
if err != nil {
- return fmt.Errorf("pod status for %s from cache failed: %v", podName, err)
+ return fmt.Errorf("pod status for %s from cache failed: %v", format.Pod(pod), err)
}
e.podLastSyncTime.Delete(pod.GetUID())
@@ -1133,9 +1102,9 @@ func (e *edged) consumePodDeletion(namespacedName *types.NamespacedName) error {
if err == apis.ErrContainerNotFound {
return err
}
- return fmt.Errorf("consume removed pod [%s] failed, %v", podName, err)
+ return fmt.Errorf("consume removed pod [%s] failed, %v", format.Pod(pod), err)
}
- klog.Infof("consume removed pod [%s] successfully\n", podName)
+ klog.Infof("consume removed pod [%s] successfully\n", format.Pod(pod))
return nil
}
@@ -1324,7 +1293,7 @@ func (e *edged) handlePod(op string, content []byte) (err error) {
case model.UpdateOperation:
e.updatePod(&pod)
case model.DeleteOperation:
- if delPod, ok := e.podManager.GetPodByName(pod.Namespace, pod.Name); ok {
+ if delPod, ok := e.podManager.GetPodByUID(pod.UID); ok {
e.deletePod(delPod)
}
}
@@ -1382,29 +1351,22 @@ func (e *edged) addPod(obj interface{}) {
attrs.OtherPods = otherpods
nodeInfo := schedulercache.NewNodeInfo(pod)
e.containerManager.UpdatePluginResources(nodeInfo, attrs)
- key := types.NamespacedName{
- Namespace: pod.Namespace,
- Name: pod.Name,
- }
e.podManager.AddPod(pod)
e.probeManager.AddPod(pod)
- e.podAdditionQueue.Add(key.String())
+ e.podAdditionQueue.Add(pod.UID)
klog.Infof("success sync addition for pod [%s]", pod.Name)
}
func (e *edged) updatePod(obj interface{}) {
newPod := obj.(*v1.Pod)
klog.Infof("start update pod [%s]", newPod.Name)
- key := types.NamespacedName{
- Namespace: newPod.Namespace,
- Name: newPod.Name,
- }
+
e.podManager.UpdatePod(newPod)
e.probeManager.AddPod(newPod)
if newPod.DeletionTimestamp == nil {
- e.podAdditionQueue.Add(key.String())
+ e.podAdditionQueue.Add(newPod.UID)
} else {
- e.podDeletionQueue.Add(key.String())
+ e.podDeletionQueue.Add(newPod.UID)
}
klog.Infof("success update pod is %+v\n", newPod)
}