summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorShelley-BaoYue <baoyue2@huawei.com>2022-06-15 11:18:53 +0800
committerShelley-BaoYue <baoyue2@huawei.com>2022-09-18 15:29:52 +0800
commitc5b53f2b5af30e91c371c27bc321d27b6d248a81 (patch)
tree3d75f03db44880e5451cf18c307186578123a3b5
parentMerge pull request #4112 from Shelley-BaoYue/feature-new-edged-stage2 (diff)
downloadkubeedge-c5b53f2b5af30e91c371c27bc321d27b6d248a81.tar.gz
add node-lease resource
Signed-off-by: Shelley-BaoYue <baoyue2@huawei.com>
-rw-r--r--cloud/pkg/cloudhub/channelq/channelq.go2
-rw-r--r--cloud/pkg/edgecontroller/controller/upstream.go143
-rw-r--r--common/constants/default.go4
-rw-r--r--edge/pkg/edged/kubeclientbridge/clientset_bridge.go7
-rw-r--r--edge/pkg/edged/kubeclientbridge/typed/coordination/v1/coordination_client_bridge.go41
-rw-r--r--edge/pkg/edged/kubeclientbridge/typed/coordination/v1/lease_bridge.go53
-rw-r--r--edge/pkg/metamanager/client/lease.go114
-rw-r--r--edge/pkg/metamanager/client/metaclient.go5
-rw-r--r--edge/pkg/metamanager/client/node.go4
-rw-r--r--edge/pkg/metamanager/metaserver/kubernetes/storage/sqlite/imitator/imitator.go9
-rw-r--r--edge/pkg/metamanager/process.go12
-rw-r--r--pkg/apis/componentconfig/cloudcore/v1alpha1/default.go4
-rw-r--r--pkg/apis/componentconfig/cloudcore/v1alpha1/types.go12
-rw-r--r--staging/src/github.com/kubeedge/beehive/pkg/core/model/message.go1
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