diff options
| author | KubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com> | 2020-02-06 22:57:20 +0800 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2020-02-06 22:57:20 +0800 |
| commit | a29fb98e0db42d7a5fb9120c413c948175c2aceb (patch) | |
| tree | fd30873abde2d1bab826045087b612f28fa5b410 | |
| parent | Merge pull request #1431 from kadisi/edgesite_dir (diff) | |
| parent | update - to . in objectSync name (diff) | |
| download | kubeedge-1.2.0-beta.0.tar.gz | |
Merge pull request #1385 from fisherxu/synccontrollerv1.2.0-beta.0
Add synccontroller for reliable message delivery
| -rw-r--r-- | cloud/cmd/cloudcore/app/server.go | 2 | ||||
| -rw-r--r-- | cloud/pkg/cloudhub/channelq/channelq.go | 5 | ||||
| -rw-r--r-- | cloud/pkg/cloudhub/handler/messagehandler.go | 3 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/clusterobjectsync.go | 3 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/config/config.go | 29 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/const.go | 20 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/createfailedobject.go | 98 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/objectsync.go | 175 | ||||
| -rw-r--r-- | cloud/pkg/synccontroller/synccontroller.go | 257 | ||||
| -rw-r--r-- | pkg/apis/cloudcore/v1alpha1/default.go | 3 | ||||
| -rw-r--r-- | pkg/apis/cloudcore/v1alpha1/types.go | 9 | ||||
| -rw-r--r-- | pkg/apis/cloudcore/v1alpha1/validation/validation.go | 11 |
12 files changed, 612 insertions, 3 deletions
diff --git a/cloud/cmd/cloudcore/app/server.go b/cloud/cmd/cloudcore/app/server.go index e6baab2ab..a9b04fbd8 100644 --- a/cloud/cmd/cloudcore/app/server.go +++ b/cloud/cmd/cloudcore/app/server.go @@ -16,6 +16,7 @@ import ( "github.com/kubeedge/kubeedge/cloud/pkg/cloudhub" "github.com/kubeedge/kubeedge/cloud/pkg/devicecontroller" "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller" + "github.com/kubeedge/kubeedge/cloud/pkg/synccontroller" "github.com/kubeedge/kubeedge/pkg/apis/cloudcore/v1alpha1" "github.com/kubeedge/kubeedge/pkg/apis/cloudcore/v1alpha1/validation" "github.com/kubeedge/kubeedge/pkg/util/flag" @@ -91,4 +92,5 @@ func registerModules(c *v1alpha1.CloudCoreConfig) { cloudhub.Register(c.Modules.CloudHub, c.KubeAPIConfig) edgecontroller.Register(c.Modules.EdgeController, c.KubeAPIConfig, "", false) devicecontroller.Register(c.Modules.DeviceController, c.KubeAPIConfig) + synccontroller.Register(c.Modules.SyncController, c.KubeAPIConfig) } diff --git a/cloud/pkg/cloudhub/channelq/channelq.go b/cloud/pkg/cloudhub/channelq/channelq.go index fb51f76ce..6ca5ab1ba 100644 --- a/cloud/pkg/cloudhub/channelq/channelq.go +++ b/cloud/pkg/cloudhub/channelq/channelq.go @@ -18,6 +18,7 @@ import ( deviceconstants "github.com/kubeedge/kubeedge/cloud/pkg/devicecontroller/constants" edgeconst "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller/constants" edgemessagelayer "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller/messagelayer" + "github.com/kubeedge/kubeedge/cloud/pkg/synccontroller" commonconst "github.com/kubeedge/kubeedge/common/constants" ) @@ -89,7 +90,7 @@ func (q *ChannelMessageQueue) addListMessageToQueue(nodeID string, msg *beehiveM } func (q *ChannelMessageQueue) addMessageToQueue(nodeID string, msg *beehiveModel.Message) { - if msg.GetResourceVersion() == "" { + if msg.GetResourceVersion() == "" && !isDeleteMessage(msg) { return } @@ -124,7 +125,7 @@ func (q *ChannelMessageQueue) addMessageToQueue(nodeID string, msg *beehiveModel return } - objectSync, err := q.ObjectSyncController.ObjectSyncLister.ObjectSyncs(resourceNamespace).Get(strings.Join([]string{nodeID, resourceUID}, "-")) + objectSync, err := q.ObjectSyncController.ObjectSyncLister.ObjectSyncs(resourceNamespace).Get(synccontroller.BuildObjectSyncName(nodeID, resourceUID)) if err == nil && msg.GetResourceVersion() <= objectSync.ResourceVersion { return } diff --git a/cloud/pkg/cloudhub/handler/messagehandler.go b/cloud/pkg/cloudhub/handler/messagehandler.go index 1c88db171..d36d4a980 100644 --- a/cloud/pkg/cloudhub/handler/messagehandler.go +++ b/cloud/pkg/cloudhub/handler/messagehandler.go @@ -24,6 +24,7 @@ import ( deviceconst "github.com/kubeedge/kubeedge/cloud/pkg/devicecontroller/constants" edgeconst "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller/constants" edgemessagelayer "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller/messagelayer" + "github.com/kubeedge/kubeedge/cloud/pkg/synccontroller" "github.com/kubeedge/kubeedge/common/constants" "github.com/kubeedge/viaduct/pkg/conn" "github.com/kubeedge/viaduct/pkg/mux" @@ -460,7 +461,7 @@ func (mh *MessageHandle) saveSuccessPoint(msg *beehiveModel.Message, info *model return } - objectSyncName := strings.Join([]string{info.NodeID, resourceUID}, "-") + objectSyncName := synccontroller.BuildObjectSyncName(info.NodeID, resourceUID) if msg.GetOperation() == beehiveModel.DeleteOperation { nodeStore.Delete(msg) diff --git a/cloud/pkg/synccontroller/clusterobjectsync.go b/cloud/pkg/synccontroller/clusterobjectsync.go new file mode 100644 index 000000000..1bcbb6fb5 --- /dev/null +++ b/cloud/pkg/synccontroller/clusterobjectsync.go @@ -0,0 +1,3 @@ +package synccontroller + +// TODO: sync for cluster level objects diff --git a/cloud/pkg/synccontroller/config/config.go b/cloud/pkg/synccontroller/config/config.go new file mode 100644 index 000000000..58916bfa2 --- /dev/null +++ b/cloud/pkg/synccontroller/config/config.go @@ -0,0 +1,29 @@ +package config + +import ( + "sync" + + "github.com/kubeedge/kubeedge/pkg/apis/cloudcore/v1alpha1" + configv1alpha1 "github.com/kubeedge/kubeedge/pkg/apis/cloudcore/v1alpha1" +) + +var c Configure +var once sync.Once + +type Configure struct { + KubeAPIConfig *v1alpha1.KubeAPIConfig + SyncController *configv1alpha1.SyncController +} + +func InitConfigure(sc *configv1alpha1.SyncController, kubeAPIConfig *v1alpha1.KubeAPIConfig) { + once.Do(func() { + c = Configure{ + KubeAPIConfig: kubeAPIConfig, + SyncController: sc, + } + }) +} + +func Get() *Configure { + return &c +} diff --git a/cloud/pkg/synccontroller/const.go b/cloud/pkg/synccontroller/const.go new file mode 100644 index 000000000..0c813a2d4 --- /dev/null +++ b/cloud/pkg/synccontroller/const.go @@ -0,0 +1,20 @@ +package synccontroller + +// Service level constants +const ( + // module + SyncControllerModuleName = "synccontroller" + CloudHubControllerModuleName = "cloudhub" + + // group + SyncControllerModuleGroup = "synccontroller" + + // kind + PodKind = "Pod" + ConfigMapKind = "ConfigMap" + SecretKind = "Secret" + ServiceKind = "Service" + EndpointKind = "Endpoints" + NodeKind = "Node" + DeviceKind = "Device" +) diff --git a/cloud/pkg/synccontroller/createfailedobject.go b/cloud/pkg/synccontroller/createfailedobject.go new file mode 100644 index 000000000..bb8c1945c --- /dev/null +++ b/cloud/pkg/synccontroller/createfailedobject.go @@ -0,0 +1,98 @@ +package synccontroller + +import ( + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/klog" + + beehiveContext "github.com/kubeedge/beehive/pkg/core/context" + "github.com/kubeedge/beehive/pkg/core/model" + edgemgr "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller/manager" + commonconst "github.com/kubeedge/kubeedge/common/constants" +) + +// Compare the objects in K8s with the objects that have been persisted to the edge, +// If the objects fail to persisted to the edge for the first time, it will be recreated here. +// Just focus on the pod, service, endpoint, because if the pod is passed down, and the configmap and secret +// cannot be found then it will be queried to the cloud. +func (sctl *SyncController) manageCreateFailedObject() { + sctl.manageCreateFailedCoreObject() + // TODO: checking for devices + // sctl.manageCreateFailedDevice() +} + +func (sctl *SyncController) manageCreateFailedCoreObject() { + allPods, err := sctl.podLister.List(labels.Everything()) + if err != nil { + klog.Errorf("Filed to list all the pods: %v", err) + } + + set := labels.Set{edgemgr.NodeRoleKey: edgemgr.NodeRoleValue} + selector := labels.SelectorFromSet(set) + allEdgeNodes, err := sctl.nodeLister.List(selector) + if err != nil { + klog.Errorf("Filed to list all the edge nodes: %v", err) + } + + for _, pod := range allPods { + if !isFromEdgeNode(allEdgeNodes, pod.Spec.NodeName) { + continue + } + // Check whether the pod is successfully persisted to edge + _, err := sctl.objectSyncLister.ObjectSyncs(pod.Namespace).Get(BuildObjectSyncName(pod.Spec.NodeName, string(pod.UID))) + if err != nil && apierrors.IsNotFound(err) { + msg := buildEdgeControllerMessage(pod.Spec.NodeName, pod.Namespace, model.ResourceTypePod, pod.Name, model.InsertOperation, pod) + beehiveContext.Send(commonconst.DefaultContextSendModuleName, *msg) + } + + // TODO: add send check for service and endpoint + /* + services, err := sctl.serviceLister.GetPodServices(pod) + if err != nil { + klog.Errorf("Filed to list all the services for pod %s: %v", pod.Name, err) + } + + for _, svc := range services { + // Check whether the service related to pod is successfully persisted to edge + _, err := sctl.objectSyncLister.ObjectSyncs(svc.Namespace).Get(buildObjectSyncName(pod.Spec.NodeName, string(svc.UID))) + if err != nil && apierrors.IsNotFound(err) { + msg := buildEdgeControllerMessage(pod.Spec.NodeName, svc.Namespace, commonconst.ResourceTypeService, svc.Name, model.InsertOperation, svc) + beehiveContext.Send(commonconst.DefaultContextSendModuleName, *msg) + } + + selector := labels.SelectorFromSet(svc.Spec.Selector) + endpoints, err := sctl.endpointLister.Endpoints(svc.Namespace).List(selector) + if err != nil { + klog.Errorf("Filed to list all the endpoints for svc %s: %v", svc.Name, err) + } + + for _, endpoint := range endpoints { + // Check whether the endpoint related to service is successfully persisted to edge + _, err := sctl.objectSyncLister.ObjectSyncs(endpoint.Namespace).Get(buildObjectSyncName(pod.Spec.NodeName, string(endpoint.UID))) + if err != nil && apierrors.IsNotFound(err) { + msg := buildEdgeControllerMessage(pod.Spec.NodeName, endpoint.Namespace, commonconst.ResourceTypeEndpoints, endpoint.Name, model.InsertOperation, endpoint) + beehiveContext.Send(commonconst.DefaultContextSendModuleName, *msg) + } + } + } + */ + } +} + +func (sctl *SyncController) manageCreateFailedDevice() { + allDevices, err := sctl.deviceLister.List(labels.Everything()) + if err != nil { + klog.Errorf("Filed to list all the devices: %v", err) + } + + for _, device := range allDevices { + // Check whether the device is successfully persisted to edge + // TODO: refactor the nodeselector of the device + nodeName := device.Spec.NodeSelector.NodeSelectorTerms[0].MatchExpressions[0].Values[0] + _, err := sctl.objectSyncLister.ObjectSyncs(device.Namespace).Get(BuildObjectSyncName(nodeName, string(device.UID))) + if err != nil && apierrors.IsNotFound(err) { + msg := buildEdgeControllerMessage(nodeName, device.Namespace, commonconst.ResourceTypeService, device.Name, model.InsertOperation, device) + beehiveContext.Send(commonconst.DefaultContextSendModuleName, *msg) + } + } +} diff --git a/cloud/pkg/synccontroller/objectsync.go b/cloud/pkg/synccontroller/objectsync.go new file mode 100644 index 000000000..e73c37ff9 --- /dev/null +++ b/cloud/pkg/synccontroller/objectsync.go @@ -0,0 +1,175 @@ +package synccontroller + +import ( + "strconv" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" + "k8s.io/klog" + + beehiveContext "github.com/kubeedge/beehive/pkg/core/context" + "github.com/kubeedge/beehive/pkg/core/model" + "github.com/kubeedge/kubeedge/cloud/pkg/apis/reliablesyncs/v1alpha1" + edgectrconst "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller/constants" + edgectrmessagelayer "github.com/kubeedge/kubeedge/cloud/pkg/edgecontroller/messagelayer" + commonconst "github.com/kubeedge/kubeedge/common/constants" +) + +func (sctl *SyncController) managePod(sync *v1alpha1.ObjectSync) { + pod, err := sctl.podLister.Pods(sync.Namespace).Get(sync.Spec.ObjectName) + + nodeName := getNodeName(sync.Name) + + if err != nil && apierrors.IsNotFound(err) { + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypePod, + "", "", nil) + return + } + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypePod, + pod.ResourceVersion, sync.Status.ObjectResourceVersion, pod) +} + +func (sctl *SyncController) manageConfigMap(sync *v1alpha1.ObjectSync) { + configmap, err := sctl.configMapLister.ConfigMaps(sync.Namespace).Get(sync.Spec.ObjectName) + + nodeName := getNodeName(sync.Name) + + if err != nil && apierrors.IsNotFound(err) { + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeConfigmap, + "", "", nil) + return + } + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeConfigmap, + configmap.ResourceVersion, sync.Status.ObjectResourceVersion, configmap) +} +func (sctl *SyncController) manageSecret(sync *v1alpha1.ObjectSync) { + secret, err := sctl.secretLister.Secrets(sync.Namespace).Get(sync.Spec.ObjectName) + + nodeName := getNodeName(sync.Name) + + if err != nil && apierrors.IsNotFound(err) { + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeSecret, + "", "", nil) + return + } + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, model.ResourceTypeSecret, + secret.ResourceVersion, sync.Status.ObjectResourceVersion, secret) +} + +func (sctl *SyncController) manageService(sync *v1alpha1.ObjectSync) { + service, err := sctl.serviceLister.Services(sync.Namespace).Get(sync.Spec.ObjectName) + + nodeName := getNodeName(sync.Name) + + if err != nil && apierrors.IsNotFound(err) { + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeService, + "", "", nil) + return + } + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeService, + service.ResourceVersion, sync.Status.ObjectResourceVersion, service) +} + +func (sctl *SyncController) manageEndpoint(sync *v1alpha1.ObjectSync) { + endpoint, err := sctl.endpointLister.Endpoints(sync.Namespace).Get(sync.Spec.ObjectName) + + nodeName := getNodeName(sync.Name) + + if err != nil && apierrors.IsNotFound(err) { + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeEndpoints, + "", "", nil) + return + } + sendEvents(err, nodeName, sync.Namespace, sync.Spec.ObjectName, commonconst.ResourceTypeEndpoints, + endpoint.ResourceVersion, sync.Status.ObjectResourceVersion, endpoint) +} + +// todo: add events for devices +func (sctl *SyncController) manageDevice(sync *v1alpha1.ObjectSync) { + //pod, err := sctl.deviceLister.Devices(sync.Namespace).Get(sync.Spec.ObjectName) + + // + //if err != nil && apierrors.IsNotFound(err) { + //trigger the delete event + //} + + //if pod.ResourceVersion > sync.Status.ObjectResourceVersion { + // trigger the update event + //} +} + +func sendEvents(err error, nodeName, namespace, objectName, resourceType string, + objectResourceVersion, syncResourceVersion string, + obj interface{}) { + + 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) + beehiveContext.Send(commonconst.DefaultContextSendModuleName, *msg) + return + } + + if CompareResourceVersion(objectResourceVersion, syncResourceVersion) > 0 { + // trigger the update event + klog.Infof("The resourceVersion: %s of %s in K8s is greater than in edgenode: %s, send the update event", objectResourceVersion, resourceType, syncResourceVersion) + msg := buildEdgeControllerMessage(nodeName, namespace, model.ResourceTypePod, objectName, model.UpdateOperation, obj) + beehiveContext.Send(commonconst.DefaultContextSendModuleName, *msg) + } +} + +func buildEdgeControllerMessage(nodeName, namespace, resourceType, resourceName, operationType string, obj interface{}) *model.Message { + msg := model.NewMessage("") + resource, err := edgectrmessagelayer.BuildResource(nodeName, namespace, resourceType, resourceName) + if err != nil { + klog.Warningf("build message resource failed with error: %s", err) + return nil + } + msg.BuildRouter(edgectrconst.EdgeControllerModuleName, edgectrconst.GroupResource, resource, operationType) + msg.Content = obj + + resourceVersion := GetObjectResourceVersion(obj) + msg.SetResourceVersion(resourceVersion) + + return msg +} + +// GetMessageUID returns the resourceVersion of the object in message +func GetObjectResourceVersion(obj interface{}) string { + if obj == nil { + klog.Error("object is nil") + return "" + } + + accessor, err := meta.Accessor(obj) + if err != nil { + klog.Errorf("Failed to get resourceVersion of the object: %v", obj) + return "" + } + + return accessor.GetResourceVersion() +} + +// CompareResourceVersion compares resourceversions, resource versions are actually +// ints, so we can easily compare them. +// If rva>rvb, return 1; rva=rvb, return 0; rva<rvb, return -1 +func CompareResourceVersion(rva, rvb string) int { + a, err := strconv.ParseUint(rva, 10, 64) + if err != nil { + // coder error + panic(err) + } + b, err := strconv.ParseUint(rvb, 10, 64) + if err != nil { + // coder error + panic(err) + } + + if a > b { + return 1 + } + if a == b { + return 0 + } + return -1 +} diff --git a/cloud/pkg/synccontroller/synccontroller.go b/cloud/pkg/synccontroller/synccontroller.go new file mode 100644 index 000000000..c1835015c --- /dev/null +++ b/cloud/pkg/synccontroller/synccontroller.go @@ -0,0 +1,257 @@ +package synccontroller + +import ( + "os" + "strings" + "time" + + "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/apimachinery/pkg/util/wait" + "k8s.io/client-go/informers" + coreinformers "k8s.io/client-go/informers/core/v1" + "k8s.io/client-go/kubernetes" + corelisters "k8s.io/client-go/listers/core/v1" + "k8s.io/client-go/rest" + "k8s.io/client-go/tools/cache" + "k8s.io/client-go/tools/clientcmd" + "k8s.io/klog" + + "github.com/kubeedge/beehive/pkg/core" + beehiveContext "github.com/kubeedge/beehive/pkg/core/context" + "github.com/kubeedge/beehive/pkg/core/model" + "github.com/kubeedge/kubeedge/cloud/pkg/apis/reliablesyncs/v1alpha1" + "github.com/kubeedge/kubeedge/cloud/pkg/client/clientset/versioned" + crdinformerfactory "github.com/kubeedge/kubeedge/cloud/pkg/client/informers/externalversions" + deviceinformer "github.com/kubeedge/kubeedge/cloud/pkg/client/informers/externalversions/devices/v1alpha1" + syncinformer "github.com/kubeedge/kubeedge/cloud/pkg/client/informers/externalversions/reliablesyncs/v1alpha1" + devicelister "github.com/kubeedge/kubeedge/cloud/pkg/client/listers/devices/v1alpha1" + synclister "github.com/kubeedge/kubeedge/cloud/pkg/client/listers/reliablesyncs/v1alpha1" + "github.com/kubeedge/kubeedge/cloud/pkg/synccontroller/config" + commonconst "github.com/kubeedge/kubeedge/common/constants" + configv1alpha1 "github.com/kubeedge/kubeedge/pkg/apis/cloudcore/v1alpha1" +) + +// SyncController use beehive context message layer +type SyncController struct { + enable bool + + // informer + podInformer coreinformers.PodInformer + configMapInformer coreinformers.ConfigMapInformer + secretInformer coreinformers.SecretInformer + serviceInformer coreinformers.ServiceInformer + endpointInformer coreinformers.EndpointsInformer + nodeInformer coreinformers.NodeInformer + deviceInformer deviceinformer.DeviceInformer + clusterObjectSyncInformer syncinformer.ClusterObjectSyncInformer + objectSyncInformer syncinformer.ObjectSyncInformer + + // synced + podSynced cache.InformerSynced + configMapSynced cache.InformerSynced + secretSynced cache.InformerSynced + serviceSynced cache.InformerSynced + endpointSynced cache.InformerSynced + nodeSynced cache.InformerSynced + deviceSynced cache.InformerSynced + clusterObjectSyncSynced cache.InformerSynced + objectSyncSynced cache.InformerSynced + + // lister + podLister corelisters.PodLister + configMapLister corelisters.ConfigMapLister + secretLister corelisters.SecretLister + serviceLister corelisters.ServiceLister + endpointLister corelisters.EndpointsLister + nodeLister corelisters.NodeLister + deviceLister devicelister.DeviceLister + clusterObjectSyncLister synclister.ClusterObjectSyncLister + objectSyncLister synclister.ObjectSyncLister +} + +func newSyncController(enable bool) *SyncController { + config, err := buildConfig() + if err != nil { + klog.Errorf("Failed to build config, err: %v", err) + os.Exit(1) + } + kubeClient := kubernetes.NewForConfigOrDie(config) + crdClient := versioned.NewForConfigOrDie(config) + + kubeSharedInformers := informers.NewSharedInformerFactory(kubeClient, 0) + crdFactory := crdinformerfactory.NewSharedInformerFactory(crdClient, 0) + + podInformer := kubeSharedInformers.Core().V1().Pods() + configMapInformer := kubeSharedInformers.Core().V1().ConfigMaps() + secretInformer := kubeSharedInformers.Core().V1().Secrets() + serviceInformer := kubeSharedInformers.Core().V1().Services() + endpointInformer := kubeSharedInformers.Core().V1().Endpoints() + nodeInformer := kubeSharedInformers.Core().V1().Nodes() + deviceInformer := crdFactory.Devices().V1alpha1().Devices() + clusterObjectSyncInformer := crdFactory.Reliablesyncs().V1alpha1().ClusterObjectSyncs() + objectSyncInformer := crdFactory.Reliablesyncs().V1alpha1().ObjectSyncs() + + sctl := &SyncController{ + enable: enable, + + podInformer: podInformer, + configMapInformer: configMapInformer, + secretInformer: secretInformer, + serviceInformer: serviceInformer, + endpointInformer: endpointInformer, + nodeInformer: nodeInformer, + deviceInformer: deviceInformer, + clusterObjectSyncInformer: clusterObjectSyncInformer, + objectSyncInformer: objectSyncInformer, + + podSynced: podInformer.Informer().HasSynced, + configMapSynced: configMapInformer.Informer().HasSynced, + secretSynced: secretInformer.Informer().HasSynced, + serviceSynced: serviceInformer.Informer().HasSynced, + endpointSynced: endpointInformer.Informer().HasSynced, + nodeSynced: nodeInformer.Informer().HasSynced, + deviceSynced: deviceInformer.Informer().HasSynced, + clusterObjectSyncSynced: clusterObjectSyncInformer.Informer().HasSynced, + objectSyncSynced: objectSyncInformer.Informer().HasSynced, + + podLister: podInformer.Lister(), + configMapLister: configMapInformer.Lister(), + secretLister: secretInformer.Lister(), + serviceLister: serviceInformer.Lister(), + endpointLister: endpointInformer.Lister(), + nodeLister: nodeInformer.Lister(), + clusterObjectSyncLister: clusterObjectSyncInformer.Lister(), + objectSyncLister: objectSyncInformer.Lister(), + } + + return sctl +} + +func Register(ec *configv1alpha1.SyncController, kubeAPIConfig *configv1alpha1.KubeAPIConfig) { + config.InitConfigure(ec, kubeAPIConfig) + core.Register(newSyncController(ec.Enable)) +} + +// Name of controller +func (sctl *SyncController) Name() string { + return SyncControllerModuleName +} + +// Group of controller +func (sctl *SyncController) Group() string { + return SyncControllerModuleGroup +} + +// Group of controller +func (sctl *SyncController) Enable() bool { + return sctl.enable +} + +// Start controller +func (sctl *SyncController) Start() { + go sctl.podInformer.Informer().Run(beehiveContext.Done()) + go sctl.configMapInformer.Informer().Run(beehiveContext.Done()) + go sctl.secretInformer.Informer().Run(beehiveContext.Done()) + go sctl.serviceInformer.Informer().Run(beehiveContext.Done()) + go sctl.endpointInformer.Informer().Run(beehiveContext.Done()) + go sctl.nodeInformer.Informer().Run(beehiveContext.Done()) + + go sctl.deviceInformer.Informer().Run(beehiveContext.Done()) + go sctl.clusterObjectSyncInformer.Informer().Run(beehiveContext.Done()) + go sctl.objectSyncInformer.Informer().Run(beehiveContext.Done()) + + if !cache.WaitForCacheSync(beehiveContext.Done(), + sctl.podSynced, + sctl.configMapSynced, + sctl.secretSynced, + sctl.serviceSynced, + sctl.endpointSynced, + sctl.nodeSynced, + sctl.deviceSynced, + sctl.clusterObjectSyncSynced, + sctl.objectSyncSynced, + ) { + klog.Errorf("unable to sync caches for sync controller") + return + } + + go wait.Until(sctl.reconcile, 5*time.Second, beehiveContext.Done()) +} + +func (sctl *SyncController) reconcile() { + allClusterObjectSyncs, err := sctl.clusterObjectSyncLister.List(labels.Everything()) + if err != nil { + klog.Errorf("Filed to list all the ClusterObjectSyncs: %v", err) + } + sctl.manageClusterObjectSync(allClusterObjectSyncs) + + allObjectSyncs, err := sctl.objectSyncLister.List(labels.Everything()) + if err != nil { + klog.Errorf("Filed to list all the ObjectSyncs: %v", err) + } + sctl.manageObjectSync(allObjectSyncs) + + sctl.manageCreateFailedObject() +} + +// Compare the cluster scope objects that have been persisted to the edge with the cluster scope objects in K8s, +// and generate update and delete events to the edge +func (sctl *SyncController) manageClusterObjectSync(syncs []*v1alpha1.ClusterObjectSync) { + // TODO: Handle cluster scope resource +} + +// Compare the namespace scope objects that have been persisted to the edge with the namespace scope objects in K8s, +// and generate update and delete events to the edge +func (sctl *SyncController) manageObjectSync(syncs []*v1alpha1.ObjectSync) { + for _, sync := range syncs { + switch sync.Spec.ObjectKind { + case model.ResourceTypePod: + sctl.managePod(sync) + case model.ResourceTypeConfigmap: + sctl.manageConfigMap(sync) + case model.ResourceTypeSecret: + sctl.manageSecret(sync) + case commonconst.ResourceTypeService: + sctl.manageService(sync) + case commonconst.ResourceTypeEndpoints: + sctl.manageEndpoint(sync) + // TODO: add device here + default: + klog.Errorf("Unsupported object kind: %v", sync.Spec.ObjectKind) + } + } +} + +// BuildObjectSyncName builds the name of objectSync/clusterObjectSync +func BuildObjectSyncName(nodeName, UID string) string { + return nodeName + "." + UID +} + +func getNodeName(syncName string) string { + tmps := strings.Split(syncName, ".") + return strings.Join(tmps[:len(tmps)-1], ".") +} + +func isFromEdgeNode(nodes []*v1.Node, nodeName string) bool { + for _, node := range nodes { + if node.Name == nodeName { + return true + } + } + return false +} + +// build Config from flags +func buildConfig() (conf *rest.Config, err error) { + kubeConfig, err := clientcmd.BuildConfigFromFlags(config.Get().KubeAPIConfig.Master, + config.Get().KubeAPIConfig.KubeConfig) + if err != nil { + return nil, err + } + kubeConfig.QPS = float32(config.Get().KubeAPIConfig.QPS) + kubeConfig.Burst = int(config.Get().KubeAPIConfig.Burst) + kubeConfig.ContentType = "application/json" + + return kubeConfig, nil +} diff --git a/pkg/apis/cloudcore/v1alpha1/default.go b/pkg/apis/cloudcore/v1alpha1/default.go index 538aeb4be..3ba6aaf67 100644 --- a/pkg/apis/cloudcore/v1alpha1/default.go +++ b/pkg/apis/cloudcore/v1alpha1/default.go @@ -120,6 +120,9 @@ func NewDefaultCloudCoreConfig() *CloudCoreConfig { UpdateDeviceStatusWorkers: constants.DefaultUpdateDeviceStatusWorkers, }, }, + SyncController: &SyncController{ + Enable: true, + }, }, } } diff --git a/pkg/apis/cloudcore/v1alpha1/types.go b/pkg/apis/cloudcore/v1alpha1/types.go index ad1a9d04f..32c6d4b9b 100644 --- a/pkg/apis/cloudcore/v1alpha1/types.go +++ b/pkg/apis/cloudcore/v1alpha1/types.go @@ -63,6 +63,8 @@ type Modules struct { EdgeController *EdgeController `json:"edgecontroller,omitempty"` // DeviceController indicates devicecontroller module config DeviceController *DeviceController `json:"devicecontroller,omitempty"` + // SyncController indicates synccontroller module config + SyncController *SyncController `json:"devicecontroller,omitempty"` } // CloudHub indicates the config of cloudhub module. @@ -294,3 +296,10 @@ type DeviceControllerLoad struct { // default 1 UpdateDeviceStatusWorkers int32 `json:"updateDeviceStatusWorkers,omitempty"` } + +// SyncController indicates the sync controller +type SyncController struct { + // Enable indicates whether devicecontroller is enabled, if set to false (for debugging etc.), skip checking other devicecontroller configs. + // default true + Enable bool `json:"enable,omitempty"` +} diff --git a/pkg/apis/cloudcore/v1alpha1/validation/validation.go b/pkg/apis/cloudcore/v1alpha1/validation/validation.go index af0a34276..1af161539 100644 --- a/pkg/apis/cloudcore/v1alpha1/validation/validation.go +++ b/pkg/apis/cloudcore/v1alpha1/validation/validation.go @@ -34,6 +34,7 @@ func ValidateCloudCoreConfiguration(c *cloudconfig.CloudCoreConfig) field.ErrorL allErrs = append(allErrs, ValidateModuleCloudHub(*c.Modules.CloudHub)...) allErrs = append(allErrs, ValidateModuleEdgeController(*c.Modules.EdgeController)...) allErrs = append(allErrs, ValidateModuleDeviceController(*c.Modules.DeviceController)...) + allErrs = append(allErrs, ValidateModuleSyncController(*c.Modules.SyncController)...) return allErrs } @@ -113,6 +114,16 @@ func ValidateModuleDeviceController(d cloudconfig.DeviceController) field.ErrorL return allErrs } +// ValidateModuleSyncController validates `d` and returns an errorList if it is invalid +func ValidateModuleSyncController(d cloudconfig.SyncController) field.ErrorList { + if !d.Enable { + return field.ErrorList{} + } + + allErrs := field.ErrorList{} + return allErrs +} + // ValidateKubeAPIConfig validates `k` and returns an errorList if it is invalid func ValidateKubeAPIConfig(k cloudconfig.KubeAPIConfig) field.ErrorList { allErrs := field.ErrorList{} |
