summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com>2020-02-28 21:13:25 +0800
committerGitHub <noreply@github.com>2020-02-28 21:13:25 +0800
commit30c42b7f3ea2c2a12c96d8648b1c4bf8a7274b3b (patch)
treebaeffff5fb697355b9f751a5488b51e9ba87c339
parentMerge pull request #1510 from fisherxu/automated-cherry-pick-of-#1491-upstrea... (diff)
parentfix nil content of delete msg which created by syccontroller (diff)
downloadkubeedge-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.go7
-rw-r--r--cloud/pkg/synccontroller/objectsync.go90
-rw-r--r--cloud/pkg/synccontroller/synccontroller.go7
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 {