diff options
| author | KubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com> | 2020-02-28 21:13:25 +0800 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2020-02-28 21:13:25 +0800 |
| commit | 30c42b7f3ea2c2a12c96d8648b1c4bf8a7274b3b (patch) | |
| tree | baeffff5fb697355b9f751a5488b51e9ba87c339 | |
| parent | Merge pull request #1510 from fisherxu/automated-cherry-pick-of-#1491-upstrea... (diff) | |
| parent | fix nil content of delete msg which created by syccontroller (diff) | |
| download | kubeedge-30c42b7f3ea2c2a12c96d8648b1c4bf8a7274b3b.tar.gz | |
Merge pull request #1512 from fisherxu/automated-cherry-pick-of-#1504-upstream-release-1.2
Automated cherry pick of #1504: fix nil content of delete msg which created by syccontroller
| -rw-r--r-- | cloud/pkg/cloudhub/channelq/channelq.go | 7 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/objectsync.go | 90 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/synccontroller.go | 7 |
3 files changed, 80 insertions, 24 deletions
diff --git a/cloud/pkg/cloudhub/channelq/channelq.go b/cloud/pkg/cloudhub/channelq/channelq.go index 9d1b1c3de..565416903 100644 --- a/cloud/pkg/cloudhub/channelq/channelq.go +++ b/cloud/pkg/cloudhub/channelq/channelq.go @@ -191,14 +191,17 @@ func isListResource(msg *beehiveModel.Message) bool { } func isDeleteMessage(msg *beehiveModel.Message) bool { + if msg.GetOperation() == beehiveModel.DeleteOperation { + return true + } deletionTimestamp, err := GetMessageDeletionTimestamp(msg) if err != nil { klog.Errorf("fail to get message DeletionTimestamp for message: %s", msg.Header.ID) return false - } - if msg.GetOperation() == beehiveModel.DeleteOperation || deletionTimestamp != nil { + } else if deletionTimestamp != nil { return true } + return false } diff --git a/cloud/pkg/synccontroller/objectsync.go b/cloud/pkg/synccontroller/objectsync.go index 312e813e1..165af1992 100644 --- a/cloud/pkg/synccontroller/objectsync.go +++ b/cloud/pkg/synccontroller/objectsync.go @@ -3,8 +3,11 @@ package synccontroller import ( "strconv" + v1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/types" "k8s.io/klog" beehiveContext "github.com/kubeedge/beehive/pkg/core/context" @@ -20,10 +23,19 @@ func (sctl *SyncController) managePod(sync *v1alpha1.ObjectSync) { nodeName := getNodeName(sync.Name) - if err != nil && apierrors.IsNotFound(err) { - sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypePod, - "", "", nil) - return + if err != nil { + if apierrors.IsNotFound(err) { + pod = &v1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: sync.Spec.ObjectName, + Namespace: sync.Namespace, + UID: types.UID(getObjectUID(sync.Name)), + }, + } + } else { + klog.Errorf("Failed to manage pod sync of %s in namespace %s: %v", sync.Name, sync.Namespace, err) + return + } } sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypePod, pod.ResourceVersion, sync.Status.ObjectResourceVersion, pod) @@ -34,10 +46,19 @@ func (sctl *SyncController) manageConfigMap(sync *v1alpha1.ObjectSync) { nodeName := getNodeName(sync.Name) - if err != nil && apierrors.IsNotFound(err) { - sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeConfigmap, - "", "", nil) - return + if err != nil { + if apierrors.IsNotFound(err) { + configmap = &v1.ConfigMap{ + ObjectMeta: metav1.ObjectMeta{ + Name: sync.Spec.ObjectName, + Namespace: sync.Namespace, + UID: types.UID(getObjectUID(sync.Name)), + }, + } + } else { + klog.Errorf("Failed to manage configMap sync of %s in namespace %s: %v", sync.Name, sync.Namespace, err) + return + } } sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeConfigmap, configmap.ResourceVersion, sync.Status.ObjectResourceVersion, configmap) @@ -48,10 +69,19 @@ func (sctl *SyncController) manageSecret(sync *v1alpha1.ObjectSync) { nodeName := getNodeName(sync.Name) - if err != nil && apierrors.IsNotFound(err) { - sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeSecret, - "", "", nil) - return + if err != nil { + if apierrors.IsNotFound(err) { + secret = &v1.Secret{ + ObjectMeta: metav1.ObjectMeta{ + Name: sync.Spec.ObjectName, + Namespace: sync.Namespace, + UID: types.UID(getObjectUID(sync.Name)), + }, + } + } else { + klog.Errorf("Failed to manage secret sync of %s in namespace %s: %v", sync.Name, sync.Namespace, err) + return + } } sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeSecret, secret.ResourceVersion, sync.Status.ObjectResourceVersion, secret) @@ -62,10 +92,19 @@ func (sctl *SyncController) manageService(sync *v1alpha1.ObjectSync) { nodeName := getNodeName(sync.Name) - if err != nil && apierrors.IsNotFound(err) { - sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeService, - "", "", nil) - return + if err != nil { + if apierrors.IsNotFound(err) { + service = &v1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: sync.Spec.ObjectName, + Namespace: sync.Namespace, + UID: types.UID(getObjectUID(sync.Name)), + }, + } + } else { + klog.Errorf("Failed to manage service sync of %s in namespace %s: %v", sync.Name, sync.Namespace, err) + return + } } sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeService, service.ResourceVersion, sync.Status.ObjectResourceVersion, service) @@ -76,10 +115,19 @@ func (sctl *SyncController) manageEndpoint(sync *v1alpha1.ObjectSync) { nodeName := getNodeName(sync.Name) - if err != nil && apierrors.IsNotFound(err) { - sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeEndpoints, - "", "", nil) - return + if err != nil { + if apierrors.IsNotFound(err) { + endpoint = &v1.Endpoints{ + ObjectMeta: metav1.ObjectMeta{ + Name: sync.Spec.ObjectName, + Namespace: sync.Namespace, + UID: types.UID(getObjectUID(sync.Name)), + }, + } + } else { + klog.Errorf("Failed to manage endpoint sync of %s in namespace %s: %v", sync.Name, sync.Namespace, err) + return + } } sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeEndpoints, endpoint.ResourceVersion, sync.Status.ObjectResourceVersion, endpoint) @@ -106,7 +154,7 @@ func sendEvents(err error, nodeName, namespace, objectName, resourceType string, if err != nil && apierrors.IsNotFound(err) { //trigger the delete event klog.Infof("%s: %s has been deleted in K8s, send the delete event to edge", resourceType, objectName) - msg := buildEdgeControllerMessage(nodeName, namespace, resourceType, objectName, model.DeleteOperation, nil) + msg := buildEdgeControllerMessage(nodeName, namespace, resourceType, objectName, model.DeleteOperation, obj) beehiveContext.Send(commonconst.DefaultContextSendModuleName, *msg) return } diff --git a/cloud/pkg/synccontroller/synccontroller.go b/cloud/pkg/synccontroller/synccontroller.go index 2af9270f1..a123a1dfa 100644 --- a/cloud/pkg/synccontroller/synccontroller.go +++ b/cloud/pkg/synccontroller/synccontroller.go @@ -5,7 +5,7 @@ import ( "strings" "time" - "k8s.io/api/core/v1" + v1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/util/wait" "k8s.io/client-go/informers" @@ -233,6 +233,11 @@ func getNodeName(syncName string) string { return strings.Join(tmps[:len(tmps)-1], ".") } +func getObjectUID(syncName string) string { + tmps := strings.Split(syncName, ".") + return tmps[len(tmps)-1] +} + func isFromEdgeNode(nodes []*v1.Node, nodeName string) bool { for _, node := range nodes { if node.Name == nodeName { |
