diff options
| author | Shelley-BaoYue <baoyue2@huawei.com> | 2022-06-15 11:18:53 +0800 |
|---|---|---|
| committer | Shelley-BaoYue <baoyue2@huawei.com> | 2022-09-18 15:29:52 +0800 |
| commit | c5b53f2b5af30e91c371c27bc321d27b6d248a81 (patch) | |
| tree | 3d75f03db44880e5451cf18c307186578123a3b5 | |
| parent | Merge pull request #4112 from Shelley-BaoYue/feature-new-edged-stage2 (diff) | |
| download | kubeedge-c5b53f2b5af30e91c371c27bc321d27b6d248a81.tar.gz | |
add node-lease resource
Signed-off-by: Shelley-BaoYue <baoyue2@huawei.com>
14 files changed, 400 insertions, 11 deletions
diff --git a/cloud/pkg/cloudhub/channelq/channelq.go b/cloud/pkg/cloudhub/channelq/channelq.go index db3ee64dc..fc5dec3fc 100644 --- a/cloud/pkg/cloudhub/channelq/channelq.go +++ b/cloud/pkg/cloudhub/channelq/channelq.go @@ -201,7 +201,7 @@ func isListResource(msg *beehiveModel.Message) bool { if msg.GetSource() == modules.EdgeControllerModuleName { resourceType, _ := messagelayer.GetResourceType(*msg) - if resourceType == beehiveModel.ResourceTypeNode { + if resourceType == beehiveModel.ResourceTypeNode || resourceType == beehiveModel.ResourceTypeLease { return true } } diff --git a/cloud/pkg/edgecontroller/controller/upstream.go b/cloud/pkg/edgecontroller/controller/upstream.go index 74248250c..13eb63053 100644 --- a/cloud/pkg/edgecontroller/controller/upstream.go +++ b/cloud/pkg/edgecontroller/controller/upstream.go @@ -34,6 +34,7 @@ import ( "time" authenticationv1 "k8s.io/api/authentication/v1" + coordinationv1 "k8s.io/api/coordination/v1" v1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/api/errors" metaV1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -41,6 +42,7 @@ import ( apimachineryType "k8s.io/apimachinery/pkg/types" k8sinformer "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" + coordinationlisters "k8s.io/client-go/listers/coordination/v1" corelisters "k8s.io/client-go/listers/core/v1" "k8s.io/klog/v2" @@ -87,6 +89,12 @@ func SortInitContainerStatuses(p *v1.Pod, statuses []v1.ContainerStatus) { } } +// ObjectResp is the object that api-server response +type ObjectResp struct { + Object metaV1.Object + Err error +} + // UpstreamController subscribe messages from edge and sync to k8s api server type UpstreamController struct { kubeClient kubernetes.Interface @@ -110,6 +118,8 @@ type UpstreamController struct { updateNodeChan chan model.Message podDeleteChan chan model.Message ruleStatusChan chan model.Message + createLeaseChan chan model.Message + queryLeaseChan chan model.Message // lister podLister corelisters.PodLister @@ -118,6 +128,7 @@ type UpstreamController struct { serviceLister corelisters.ServiceLister endpointLister corelisters.EndpointsLister nodeLister corelisters.NodeLister + leaseLister coordinationlisters.LeaseLister } // Start UpstreamController @@ -159,6 +170,12 @@ func (uc *UpstreamController) Start() error { for i := 0; i < int(uc.config.Load.DeletePodWorkers); i++ { go uc.deletePod() } + for i := 0; i < int(uc.config.Load.CreateLeaseWorkers); i++ { + go uc.createOrUpdateLease() + } + for i := 0; i < int(uc.config.Load.QueryLeaseWorkers); i++ { + go uc.queryLease() + } for i := 0; i < int(uc.config.Load.UpdateRuleStatusWorkers); i++ { go uc.updateRuleStatus() } @@ -224,12 +241,19 @@ func (uc *UpstreamController) dispatchMessage() { } case model.ResourceTypeRuleStatus: uc.ruleStatusChan <- msg - + case model.ResourceTypeLease: + switch msg.GetOperation() { + case model.InsertOperation, model.UpdateOperation: + uc.createLeaseChan <- msg + case model.QueryOperation: + uc.queryLeaseChan <- msg + } default: klog.Errorf("message: %s, resource type: %s unsupported", msg.GetID(), resourceType) } } } + func (uc *UpstreamController) updateRuleStatus() { for { select { @@ -589,6 +613,8 @@ func kubeClientGet(uc *UpstreamController, namespace string, name string, queryT obj, err = uc.nodeLister.Get(name) case model.ResourceTypeServiceAccountToken: obj, err = uc.getServiceAccountToken(namespace, name, msg) + case model.ResourceTypeLease: + obj, err = uc.leaseLister.Leases(namespace).Get(name) default: err = stderrors.New("wrong query type") } @@ -895,6 +921,118 @@ func (uc *UpstreamController) queryNode() { } } +func (uc *UpstreamController) createOrUpdateLease() { + for { + select { + case <-beehiveContext.Done(): + klog.Warning("stop create or update lease") + return + case msg := <-uc.createLeaseChan: + klog.V(4).Infof("message: %s, operation is: %s, and resource is: %s", msg.GetID(), msg.GetOperation(), msg.GetResource()) + + data, err := msg.GetContentData() + if err != nil { + klog.Warningf("message: %s process failure, get content data failed with error: %v", msg.GetID(), err) + continue + } + + namespace, err := messagelayer.GetNamespace(msg) + if err != nil { + klog.Warningf("message: %s process failure, get namespace failed with error: %v", msg.GetID(), err) + return + } + name, err := messagelayer.GetResourceName(msg) + if err != nil { + klog.Warningf("message: %s process failure, get resource name failed with error: %v", msg.GetID(), err) + continue + } + + lease := &coordinationv1.Lease{} + err = json.Unmarshal(data, lease) + if err != nil { + errLog := fmt.Sprintf("message: %s process failure, unmarshal message content with error: %v", msg.GetID(), err) + klog.Error(errLog) + uc.nodeMsgResponse(name, namespace, errLog, msg) + continue + } + + switch msg.GetOperation() { + case model.InsertOperation: + resp, err := uc.kubeClient.CoordinationV1().Leases(namespace).Create(context.TODO(), lease, metaV1.CreateOptions{}) + if err != nil { + klog.Errorf("create lease %s failed, error: %v", name, err) + } + + resMsg := model.NewMessage(msg.GetID()). + FillBody(&ObjectResp{Object: resp, Err: err}). + BuildRouter(modules.EdgeControllerModuleName, constants.GroupResource, msg.GetResource(), model.ResponseOperation) + if err = uc.messageLayer.Response(*resMsg); err != nil { + klog.Warningf("Response message: %s failed, response failed with error: %v", msg.GetID(), err) + continue + } + + klog.V(4).Infof("message: %s, create lease successfully, namespace: %s, name: %s", msg.GetID(), namespace, name) + + case model.UpdateOperation: + resp, err := uc.kubeClient.CoordinationV1().Leases(namespace).Update(context.TODO(), lease, metaV1.UpdateOptions{}) + if err != nil { + klog.Error("Update lease %s failed, error: %v", name, err) + } + + resMsg := model.NewMessage(msg.GetID()). + FillBody(&ObjectResp{Object: resp, Err: err}). + BuildRouter(modules.EdgeControllerModuleName, constants.GroupResource, msg.GetResource(), model.ResponseOperation) + if err = uc.messageLayer.Response(*resMsg); err != nil { + klog.Warningf("Response message: %s failed, response failed with error: %v", msg.GetID(), err) + continue + } + + klog.V(4).Infof("message: %s, update lease successfully, namespace: %s, name: %s", msg.GetID(), namespace, name) + + default: + klog.Warningf("message: %s process failure, operation: %s unsupported", msg.GetID(), msg.GetOperation()) + } + } + } +} + +func (uc *UpstreamController) queryLease() { + for { + select { + case <-beehiveContext.Done(): + klog.Warning("stop queryLease") + return + case msg := <-uc.queryLeaseChan: + klog.V(4).Infof("message: %s, operation is: %s, and resource is: %s", msg.GetID(), msg.GetOperation(), msg.GetResource()) + namespace, err := messagelayer.GetNamespace(msg) + if err != nil { + klog.Warningf("message: %s process failure, get namespace failed with error: %v", msg.GetID(), err) + return + } + name, err := messagelayer.GetResourceName(msg) + if err != nil { + klog.Warningf("message: %s process failure, get resource name failed with error: %v", msg.GetID(), err) + return + } + + object, err := kubeClientGet(uc, namespace, name, model.ResourceTypeLease, msg) + if err != nil { + klog.Error("Query lease %s failed, error: %v", name, err) + } + + resMsg := model.NewMessage(msg.GetID()). + FillBody(&ObjectResp{Object: object, Err: err}). + BuildRouter(modules.EdgeControllerModuleName, constants.GroupResource, msg.GetResource(), model.ResponseOperation) + if err = uc.messageLayer.Response(*resMsg); err != nil { + klog.Warningf("Response message: %s failed, response failed with error: %v", msg.GetID(), err) + continue + } + + klog.V(4).Infof("message: %s, query lease successfully, namespace: %s, name: %s", msg.GetID(), namespace, name) + } + } +} + func (uc *UpstreamController) unmarshalPodStatusMessage(msg model.Message) (ns string, podStatuses []edgeapi.PodStatusRequest) { ns, err := messagelayer.GetNamespace(msg) if err != nil { @@ -1030,6 +1168,7 @@ func NewUpstreamController(config *v1alpha1.EdgeController, factory k8sinformer. uc.podLister = factory.Core().V1().Pods().Lister() uc.configMapLister = factory.Core().V1().ConfigMaps().Lister() uc.secretLister = factory.Core().V1().Secrets().Lister() + uc.leaseLister = factory.Coordination().V1().Leases().Lister() uc.nodeStatusChan = make(chan model.Message, config.Buffer.UpdateNodeStatus) uc.podStatusChan = make(chan model.Message, config.Buffer.UpdatePodStatus) @@ -1044,6 +1183,8 @@ func NewUpstreamController(config *v1alpha1.EdgeController, factory k8sinformer. uc.queryNodeChan = make(chan model.Message, config.Buffer.QueryNode) uc.updateNodeChan = make(chan model.Message, config.Buffer.UpdateNode) uc.podDeleteChan = make(chan model.Message, config.Buffer.DeletePod) + uc.createLeaseChan = make(chan model.Message, config.Buffer.CreateLease) + uc.queryLeaseChan = make(chan model.Message, config.Buffer.QueryLease) uc.ruleStatusChan = make(chan model.Message, config.Buffer.UpdateNodeStatus) return uc, nil } diff --git a/common/constants/default.go b/common/constants/default.go index 75d686849..bfec31f46 100644 --- a/common/constants/default.go +++ b/common/constants/default.go @@ -91,6 +91,8 @@ const ( DefaultUpdateNodeWorkers = 4 DefaultDeletePodWorkers = 4 DefaultUpdateRuleStatusWorkers = 4 + DefaultCreateLeaseWorkers = 4 + DefaultQueryLeaseWorkers = 4 DefaultServiceAccountTokenWorkers = 4 DefaultUpdatePodStatusBuffer = 1024 @@ -105,6 +107,8 @@ const ( DefaultQueryNodeBuffer = 1024 DefaultUpdateNodeBuffer = 1024 DefaultDeletePodBuffer = 1024 + DefaultCreateLeaseBuffer = 1024 + DefaultQueryLeaseBuffer = 1024 DefaultServiceAccountTokenBuffer = 1024 DefaultPodEventBuffer = 1 diff --git a/edge/pkg/edged/kubeclientbridge/clientset_bridge.go b/edge/pkg/edged/kubeclientbridge/clientset_bridge.go index 6dadd89a7..3548ee096 100644 --- a/edge/pkg/edged/kubeclientbridge/clientset_bridge.go +++ b/edge/pkg/edged/kubeclientbridge/clientset_bridge.go @@ -26,11 +26,14 @@ package kubeclientbridge import ( clientset "k8s.io/client-go/kubernetes" fakekube "k8s.io/client-go/kubernetes/fake" + coordinationv1 "k8s.io/client-go/kubernetes/typed/coordination/v1" + fakecoordinationv1 "k8s.io/client-go/kubernetes/typed/coordination/v1/fake" corev1 "k8s.io/client-go/kubernetes/typed/core/v1" fakecorev1 "k8s.io/client-go/kubernetes/typed/core/v1/fake" storagev1 "k8s.io/client-go/kubernetes/typed/storage/v1" fakestoragev1 "k8s.io/client-go/kubernetes/typed/storage/v1/fake" + kecoordinationv1 "github.com/kubeedge/kubeedge/edge/pkg/edged/kubeclientbridge/typed/coordination/v1" kecorev1 "github.com/kubeedge/kubeedge/edge/pkg/edged/kubeclientbridge/typed/core/v1" kestoragev1 "github.com/kubeedge/kubeedge/edge/pkg/edged/kubeclientbridge/typed/storage/v1" "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client" @@ -56,3 +59,7 @@ func (c *Clientset) CoreV1() corev1.CoreV1Interface { func (c *Clientset) StorageV1() storagev1.StorageV1Interface { return &kestoragev1.StorageV1Bridge{FakeStorageV1: fakestoragev1.FakeStorageV1{Fake: &c.Fake}, MetaClient: c.MetaClient} } + +func (c *Clientset) CoordinationV1() coordinationv1.CoordinationV1Interface { + return &kecoordinationv1.CoordinationV1Bridge{FakeCoordinationV1: fakecoordinationv1.FakeCoordinationV1{Fake: &c.Fake}, MetaClient: c.MetaClient} +} diff --git a/edge/pkg/edged/kubeclientbridge/typed/coordination/v1/coordination_client_bridge.go b/edge/pkg/edged/kubeclientbridge/typed/coordination/v1/coordination_client_bridge.go new file mode 100644 index 000000000..e628c1b56 --- /dev/null +++ b/edge/pkg/edged/kubeclientbridge/typed/coordination/v1/coordination_client_bridge.go @@ -0,0 +1,41 @@ +/* +Copyright 2022 The Kubernetes Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. + +@CHANGELOG +KubeEdge Authors: To make a bridge between kubeclient and metaclient, +This file is derived from K8S client-go code with reduced set of methods +Changes done are +1. Package v1 got some functions from "k8s.io/client-go/kubernetes/typed/coordination/v1/fake/fake_coordination_client.go" +and made some variant +*/ + +package v1 + +import ( + v1 "k8s.io/client-go/kubernetes/typed/coordination/v1" + fakecoordinationv1 "k8s.io/client-go/kubernetes/typed/coordination/v1/fake" + + "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client" +) + +// CoordinationV1Bridge is a coordinationV1 bridge +type CoordinationV1Bridge struct { + fakecoordinationv1.FakeCoordinationV1 + MetaClient client.CoreInterface +} + +func (c *CoordinationV1Bridge) Leases(namespace string) v1.LeaseInterface { + return &LeaseBridge{fakecoordinationv1.FakeLeases{Fake: &c.FakeCoordinationV1}, namespace, c.MetaClient} +} diff --git a/edge/pkg/edged/kubeclientbridge/typed/coordination/v1/lease_bridge.go b/edge/pkg/edged/kubeclientbridge/typed/coordination/v1/lease_bridge.go new file mode 100644 index 000000000..f2825b023 --- /dev/null +++ b/edge/pkg/edged/kubeclientbridge/typed/coordination/v1/lease_bridge.go @@ -0,0 +1,53 @@ +/* +Copyright 2022 The Kubernetes Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. + +@CHANGELOG +KubeEdge Authors: To make a bridge between kubeclient and metaclient, +This file is derived from K8S client-go code with reduced set of methods +Changes done are +1. Package v1 got some functions from "k8s.io/client-go/kubernetes/typed/coordination/v1/fake/fake_lease.go" +and made some variant +*/ + +package v1 + +import ( + "context" + + coordinationv1 "k8s.io/api/coordination/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + fakecoordinationv1 "k8s.io/client-go/kubernetes/typed/coordination/v1/fake" + + "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client" +) + +// LeaseBridge implements LeaseInterface +type LeaseBridge struct { + fakecoordinationv1.FakeLeases + ns string + MetaClient client.CoreInterface +} + +func (c *LeaseBridge) Create(ctx context.Context, lease *coordinationv1.Lease, opts metav1.CreateOptions) (result *coordinationv1.Lease, err error) { + return c.MetaClient.Leases(c.ns).Create(lease) +} + +func (c *LeaseBridge) Update(ctx context.Context, lease *coordinationv1.Lease, opts metav1.UpdateOptions) (result *coordinationv1.Lease, err error) { + return c.MetaClient.Leases(c.ns).Update(lease) +} + +func (c *LeaseBridge) Get(ctx context.Context, name string, options metav1.GetOptions) (result *coordinationv1.Lease, err error) { + return c.MetaClient.Leases(c.ns).Get(name) +} diff --git a/edge/pkg/metamanager/client/lease.go b/edge/pkg/metamanager/client/lease.go new file mode 100644 index 000000000..82ec0f898 --- /dev/null +++ b/edge/pkg/metamanager/client/lease.go @@ -0,0 +1,114 @@ +/* +Copyright 2022 The KubeEdge Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package client + +import ( + "encoding/json" + "fmt" + + coordinationv1 "k8s.io/api/coordination/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + + "github.com/kubeedge/beehive/pkg/core/model" + "github.com/kubeedge/kubeedge/edge/pkg/common/message" + "github.com/kubeedge/kubeedge/edge/pkg/common/modules" +) + +// LeasesGetter to get lease interface +type LeasesGetter interface { + Leases(namespace string) LeasesInterface +} + +// LeasesInterface is interface for client leases +type LeasesInterface interface { + Create(lease *coordinationv1.Lease) (*coordinationv1.Lease, error) + Update(lease *coordinationv1.Lease) (*coordinationv1.Lease, error) + Get(name string) (*coordinationv1.Lease, error) +} + +type leases struct { + namespace string + send SendInterface +} + +func newLeases(namespace string, s SendInterface) *leases { + return &leases{ + send: s, + namespace: namespace, + } +} + +// LeaseResp represents lease response from the api-server +type LeaseResp struct { + Object *coordinationv1.Lease + Err *apierrors.StatusError +} + +func (c *leases) Create(lease *coordinationv1.Lease) (*coordinationv1.Lease, error) { + resource := fmt.Sprintf("%s/%s/%s", c.namespace, model.ResourceTypeLease, lease.Name) + leaseMsg := message.BuildMsg(modules.MetaGroup, "", modules.EdgedModuleName, resource, model.InsertOperation, lease) + resp, err := c.send.SendSync(leaseMsg) + if err != nil { + return nil, fmt.Errorf("create lease failed, err: %v", err) + } + + content, err := resp.GetContentData() + if err != nil { + return nil, fmt.Errorf("parse message to lease failed, err: %v", err) + } + return handleLeaseResp(content) +} + +func (c *leases) Update(lease *coordinationv1.Lease) (*coordinationv1.Lease, error) { + resource := fmt.Sprintf("%s/%s/%s", c.namespace, model.ResourceTypeLease, lease.Name) + leaseMsg := message.BuildMsg(modules.MetaGroup, "", modules.EdgedModuleName, resource, model.UpdateOperation, lease) + resp, err := c.send.SendSync(leaseMsg) + if err != nil { + return nil, fmt.Errorf("update lease failed, err: %v", err) + } + + content, err := resp.GetContentData() + if err != nil { + return nil, fmt.Errorf("parse message to lease failed, err: %v", err) + } + return handleLeaseResp(content) +} + +func (c *leases) Get(name string) (*coordinationv1.Lease, error) { + resource := fmt.Sprintf("%s/%s/%s", c.namespace, model.ResourceTypeLease, name) + leaseMsg := message.BuildMsg(modules.MetaGroup, "", modules.EdgedModuleName, resource, model.QueryOperation, nil) + resp, err := c.send.SendSync(leaseMsg) + if err != nil { + return nil, fmt.Errorf("query lease failed, err: %v", err) + } + + content, err := resp.GetContentData() + if err != nil { + return nil, fmt.Errorf("parse message to lease failed, err: %v", err) + } + return handleLeaseResp(content) +} + +func handleLeaseResp(content []byte) (*coordinationv1.Lease, error) { + var leaseResp LeaseResp + err := json.Unmarshal(content, &leaseResp) + if err != nil { + return nil, fmt.Errorf("unmarshal message to lease failed, err: %v", err) + } + + return leaseResp.Object, leaseResp.Err +} diff --git a/edge/pkg/metamanager/client/metaclient.go b/edge/pkg/metamanager/client/metaclient.go index 8ef17d1e3..1b3c72106 100644 --- a/edge/pkg/metamanager/client/metaclient.go +++ b/edge/pkg/metamanager/client/metaclient.go @@ -28,6 +28,7 @@ type CoreInterface interface { PersistentVolumesGetter PersistentVolumeClaimsGetter VolumeAttachmentsGetter + LeasesGetter } type metaClient struct { @@ -77,6 +78,10 @@ func (m *metaClient) VolumeAttachments(namespace string) VolumeAttachmentsInterf return newVolumeAttachments(namespace, m.send) } +func (m *metaClient) Leases(namespace string) LeasesInterface { + return newLeases(namespace, m.send) +} + // New creates new metaclient func New() CoreInterface { return &metaClient{ diff --git a/edge/pkg/metamanager/client/node.go b/edge/pkg/metamanager/client/node.go index 9225c5256..d24b26085 100644 --- a/edge/pkg/metamanager/client/node.go +++ b/edge/pkg/metamanager/client/node.go @@ -11,12 +11,12 @@ import ( "github.com/kubeedge/kubeedge/edge/pkg/common/modules" ) -//NodesGetter to get node interface +// NodesGetter to get node interface type NodesGetter interface { Nodes(namespace string) NodesInterface } -//NodesInterface is interface for client nodes +// NodesInterface is interface for client nodes type NodesInterface interface { Create(*api.Node) (*api.Node, error) Update(*api.Node) error diff --git a/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/imitator.go b/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/imitator.go index 0d074da76..47c3cc412 100644 --- a/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/imitator.go +++ b/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/imitator.go @@ -18,6 +18,7 @@ import ( "github.com/kubeedge/beehive/pkg/core/model" "github.com/kubeedge/kubeedge/common/constants" "github.com/kubeedge/kubeedge/edge/pkg/common/dbm" + "github.com/kubeedge/kubeedge/edge/pkg/common/modules" v2 "github.com/kubeedge/kubeedge/edge/pkg/metamanager/dao/v2" "github.com/kubeedge/kubeedge/edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/watchhook" "github.com/kubeedge/kubeedge/pkg/metaserver" @@ -199,10 +200,10 @@ func (s *imitator) Event(msg *model.Message) []watch.Event { klog.V(4).Infof("[metaserver] get a message from metamanager: %+v", msg) var ret []watch.Event _, resType, _ := parseResource(msg.Router.Resource) - //skip nodestatus and podstatus - if strings.Contains(resType, "status") { - klog.V(4).Infof("skip status messages") - return ret + //skip nodestatus, podstatus and node-lease + if strings.Contains(resType, "status") || (strings.Contains(resType, "lease") && msg.GetSource() == modules.EdgedModuleName) { + klog.V(4).Infof("skip status or node-lease messages") + return []watch.Event{} } var bytes []byte var err error diff --git a/edge/pkg/metamanager/process.go b/edge/pkg/metamanager/process.go index 79f840ee3..311773838 100644 --- a/edge/pkg/metamanager/process.go +++ b/edge/pkg/metamanager/process.go @@ -84,7 +84,8 @@ func requireRemoteQuery(resType string) bool { resType == constants.ResourceTypePersistentVolumeClaim || resType == constants.ResourceTypeVolumeAttachment || resType == model.ResourceTypeNode || - resType == model.ResourceTypeServiceAccountToken + resType == model.ResourceTypeServiceAccountToken || + resType == model.ResourceTypeLease } func isConnected() bool { @@ -128,6 +129,11 @@ func (m *metaManager) processInsert(message model.Message) { return } + if resType == model.ResourceTypeLease && message.GetSource() == modules.EdgedModuleName { + sendToCloud(&message) + return + } + // Notify edged sendToEdged(&message, false) @@ -171,7 +177,7 @@ func (m *metaManager) processUpdate(message model.Message) { sendToCloud(&message) // For pod status update message, we need to wait for the response message // to ensure that the pod status is correctly reported to the kube-apiserver - if resType != model.ResourceTypePodStatus { + if resType != model.ResourceTypePodStatus && resType != model.ResourceTypeLease { resp := message.NewRespByMessage(&message, OK) sendToEdged(resp, message.IsSync()) } @@ -258,7 +264,7 @@ func (m *metaManager) processQuery(message model.Message) { var err error if requireRemoteQuery(resType) && isConnected() { metas, err = dao.QueryMeta("key", resKey) - if err != nil || len(*metas) == 0 || resType == model.ResourceTypeNode || resType == constants.ResourceTypeVolumeAttachment { + if err != nil || len(*metas) == 0 || resType == model.ResourceTypeNode || resType == constants.ResourceTypeVolumeAttachment || resType == model.ResourceTypeLease { m.processRemoteQuery(message) } else { resp := message.NewRespByMessage(&message, *metas) diff --git a/pkg/apis/componentconfig/cloudcore/v1alpha1/default.go b/pkg/apis/componentconfig/cloudcore/v1alpha1/default.go index 7f7b67f4f..ce6be8db5 100644 --- a/pkg/apis/componentconfig/cloudcore/v1alpha1/default.go +++ b/pkg/apis/componentconfig/cloudcore/v1alpha1/default.go @@ -100,6 +100,8 @@ func NewDefaultCloudCoreConfig() *CloudCoreConfig { QueryNode: constants.DefaultQueryNodeBuffer, UpdateNode: constants.DefaultUpdateNodeBuffer, DeletePod: constants.DefaultDeletePodBuffer, + CreateLease: constants.DefaultCreateLeaseBuffer, + QueryLease: constants.DefaultQueryLeaseBuffer, ServiceAccountToken: constants.DefaultServiceAccountTokenBuffer, }, Load: &EdgeControllerLoad{ @@ -115,6 +117,8 @@ func NewDefaultCloudCoreConfig() *CloudCoreConfig { QueryNodeWorkers: constants.DefaultQueryNodeWorkers, UpdateNodeWorkers: constants.DefaultUpdateNodeWorkers, DeletePodWorkers: constants.DefaultDeletePodWorkers, + CreateLeaseWorkers: constants.DefaultCreateLeaseWorkers, + QueryLeaseWorkers: constants.DefaultQueryLeaseWorkers, UpdateRuleStatusWorkers: constants.DefaultUpdateRuleStatusWorkers, ServiceAccountTokenWorkers: constants.DefaultServiceAccountTokenWorkers, }, diff --git a/pkg/apis/componentconfig/cloudcore/v1alpha1/types.go b/pkg/apis/componentconfig/cloudcore/v1alpha1/types.go index 46384410e..fd3e01ae7 100644 --- a/pkg/apis/componentconfig/cloudcore/v1alpha1/types.go +++ b/pkg/apis/componentconfig/cloudcore/v1alpha1/types.go @@ -263,6 +263,12 @@ type EdgeControllerBuffer struct { // DeletePod indicates the buffer of delete pod message from edge // default 1024 DeletePod int32 `json:"deletePod,omitempty"` + // CreateLease indicates the buffer of create lease message from edge + // default 1024 + CreateLease int32 `json:"createLease,omitempty"` + // QueryLease indicates the buffer of query lease message from edge + // default 1024 + QueryLease int32 `json:"queryLease,omitempty"` // ServiceAccount indicates the buffer of service account token // default 1024 ServiceAccountToken int32 `json:"serviceAccountToken,omitempty"` @@ -306,6 +312,12 @@ type EdgeControllerLoad struct { // DeletePodWorkers indicates the load of delete pod workers // default 4 DeletePodWorkers int32 `json:"deletePodWorkers,omitempty"` + // CreateLeaseWorkers indicates the load of create lease workers + // default 4 + CreateLeaseWorkers int32 `json:"createLeaseWorkers,omitempty"` + // QueryLeaseWorkers indicates the load of query lease workers + // default 4 + QueryLeaseWorkers int32 `json:"queryLeaseWorkers,omitempty"` // UpdateRuleStatusWorkers indicates the load of update rule status // default 4 UpdateRuleStatusWorkers int32 `json:"UpdateRuleStatusWorkers,omitempty"` diff --git a/staging/src/github.com/kubeedge/beehive/pkg/core/model/message.go b/staging/src/github.com/kubeedge/beehive/pkg/core/model/message.go index eb1ac9c8c..abbc1aedb 100644 --- a/staging/src/github.com/kubeedge/beehive/pkg/core/model/message.go +++ b/staging/src/github.com/kubeedge/beehive/pkg/core/model/message.go @@ -31,6 +31,7 @@ const ( ResourceTypeRule = "rule" ResourceTypeRuleEndpoint = "ruleendpoint" ResourceTypeRuleStatus = "rulestatus" + ResourceTypeLease = "lease" ) // Message struct |
