summaryrefslogtreecommitdiff
path: root/edge
diff options
context:
space:
mode:
authorliuzhiyi1993 <liuzhiyi@huawei.com>2020-02-13 05:56:21 +0800
committerliuzhiyi1993 <liuzhiyi@huawei.com>2020-03-31 22:30:02 +0800
commit0dfe63446ad76e41a084ade9d3c393a2554def34 (patch)
treef9857c1fd4ad3e9a00aa04effdb1c57bde6398b7 /edge
parentMerge pull request #1567 from chendave/broken (diff)
downloadkubeedge-0dfe63446ad76e41a084ade9d3c393a2554def34.tar.gz
refactor edgemesh
Diffstat (limited to 'edge')
-rw-r--r--edge/pkg/metamanager/client/listener.go78
-rw-r--r--edge/pkg/metamanager/client/metaclient.go6
-rw-r--r--edge/pkg/metamanager/client/service.go71
-rw-r--r--edge/pkg/metamanager/process.go74
4 files changed, 206 insertions, 23 deletions
diff --git a/edge/pkg/metamanager/client/listener.go b/edge/pkg/metamanager/client/listener.go
new file mode 100644
index 000000000..046df3c59
--- /dev/null
+++ b/edge/pkg/metamanager/client/listener.go
@@ -0,0 +1,78 @@
+package client
+
+import (
+ "fmt"
+
+ "github.com/kubeedge/beehive/pkg/core/model"
+
+ "github.com/kubeedge/kubeedge/common/constants"
+ "github.com/kubeedge/kubeedge/edge/pkg/common/message"
+ "github.com/kubeedge/kubeedge/edge/pkg/common/modules"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/constant"
+)
+
+// Listener is only used by EdgeMesh. It stores
+// the fakeIP of EdgeMesh into edge db. One fakeIP for
+// one service.
+const (
+ DefaultNamespace = "default"
+)
+
+// ListenerGetter interface
+type ListenerGetter interface {
+ Listener() ListenInterface
+}
+
+// ListenInterface is an interface
+type ListenInterface interface {
+ Add(key interface{}, value interface{}) error
+ Del(key interface{}) error
+ Get(key interface{}) (interface{}, error)
+}
+
+type listener struct {
+ send SendInterface
+}
+
+func newListener(s SendInterface) *listener {
+ return &listener{
+ send: s,
+ }
+}
+
+func (ln *listener) Add(key interface{}, value interface{}) error {
+ svcName, ok := key.(string)
+ if !ok {
+ return fmt.Errorf("the key type is invalid")
+ }
+ resource := fmt.Sprintf("%s/%s/%s", DefaultNamespace, constants.ResourceTypeListener, svcName)
+ msg := message.BuildMsg(modules.MetaGroup, "", constant.ModuleNameEdgeMesh, resource, model.InsertOperation, value)
+ _, err := ln.send.SendSync(msg)
+ return err
+}
+
+func (ln *listener) Del(key interface{}) error {
+ svcName, ok := key.(string)
+ if !ok {
+ return fmt.Errorf("the key type is invalid")
+ }
+ resource := fmt.Sprintf("%s/%s/%s", DefaultNamespace, constants.ResourceTypeListener, svcName)
+ msg := message.BuildMsg(modules.MetaGroup, "", constant.ModuleNameEdgeMesh, resource, model.DeleteOperation, nil)
+ _, err := ln.send.SendSync(msg)
+ return err
+}
+
+func (ln *listener) Get(key interface{}) (interface{}, error) {
+ svcName, ok := key.(string)
+ if !ok {
+ return nil, fmt.Errorf("the key type is invalid")
+ }
+ resource := fmt.Sprintf("%s/%s/%s", DefaultNamespace, constants.ResourceTypeListener, svcName)
+ msg := message.BuildMsg(modules.MetaGroup, "", constant.ModuleNameEdgeMesh, resource, model.QueryOperation, nil)
+ respMsg, err := ln.send.SendSync(msg)
+ if err != nil {
+ return nil, err
+ }
+
+ return respMsg.Content, nil
+}
diff --git a/edge/pkg/metamanager/client/metaclient.go b/edge/pkg/metamanager/client/metaclient.go
index c5d81950b..5d45ee0cc 100644
--- a/edge/pkg/metamanager/client/metaclient.go
+++ b/edge/pkg/metamanager/client/metaclient.go
@@ -29,6 +29,7 @@ type CoreInterface interface {
PersistentVolumesGetter
PersistentVolumeClaimsGetter
VolumeAttachmentsGetter
+ ListenerGetter
}
type metaClient struct {
@@ -84,6 +85,11 @@ func (m *metaClient) VolumeAttachments(namespace string) VolumeAttachmentsInterf
return newVolumeAttachments(namespace, m.send)
}
+// New Listener metaClient
+func (m *metaClient) Listener() ListenInterface {
+ return newListener(m.send)
+}
+
// New creates new metaclient
func New() CoreInterface {
return &metaClient{
diff --git a/edge/pkg/metamanager/client/service.go b/edge/pkg/metamanager/client/service.go
index f60b440e9..edfb16f56 100644
--- a/edge/pkg/metamanager/client/service.go
+++ b/edge/pkg/metamanager/client/service.go
@@ -4,9 +4,9 @@ import (
"encoding/json"
"fmt"
+ "github.com/kubeedge/beehive/pkg/core/model"
v1 "k8s.io/api/core/v1"
- "github.com/kubeedge/beehive/pkg/core/model"
"github.com/kubeedge/kubeedge/common/constants"
"github.com/kubeedge/kubeedge/edge/pkg/common/message"
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
@@ -25,6 +25,7 @@ type ServiceInterface interface {
Delete(name string) error
Get(name string) (*v1.Service, error)
GetPods(name string) ([]v1.Pod, error)
+ ListAll() ([]v1.Service, error)
}
type services struct {
@@ -70,12 +71,12 @@ func (s *services) GetPods(name string) ([]v1.Pod, error) {
}
if respMsg.GetOperation() == model.ResponseOperation {
- return handlerServicePodListFromMetaDB(content)
+ return handleServicePodListFromMetaDB(content)
}
- return handlerServicePodListFromMetaManager(content)
+ return handleServicePodListFromMetaManager(content)
}
-func handlerServicePodListFromMetaDB(content []byte) ([]v1.Pod, error) {
+func handleServicePodListFromMetaDB(content []byte) ([]v1.Pod, error) {
var lists []string
err := json.Unmarshal([]byte(content), &lists)
if err != nil {
@@ -83,7 +84,7 @@ func handlerServicePodListFromMetaDB(content []byte) ([]v1.Pod, error) {
}
if len(lists) != 1 {
- return nil, fmt.Errorf("Service length from meta db is %d", len(lists))
+ return nil, fmt.Errorf("service length from meta db is %d", len(lists))
}
var ps []v1.Pod
@@ -94,7 +95,7 @@ func handlerServicePodListFromMetaDB(content []byte) ([]v1.Pod, error) {
return ps, nil
}
-func handlerServicePodListFromMetaManager(content []byte) ([]v1.Pod, error) {
+func handleServicePodListFromMetaManager(content []byte) ([]v1.Pod, error) {
var ps []v1.Pod
err := json.Unmarshal(content, &ps)
if err != nil {
@@ -122,12 +123,12 @@ func (s *services) Get(name string) (*v1.Service, error) {
}
if respMsg.GetOperation() == model.ResponseOperation {
- return handlerServiceFromMetaDB(content)
+ return handleServiceFromMetaDB(content)
}
return handleServiceFromMetaManager(content)
}
-func handlerServiceFromMetaDB(content []byte) (*v1.Service, error) {
+func handleServiceFromMetaDB(content []byte) (*v1.Service, error) {
var lists []string
err := json.Unmarshal([]byte(content), &lists)
if err != nil {
@@ -135,7 +136,7 @@ func handlerServiceFromMetaDB(content []byte) (*v1.Service, error) {
}
if len(lists) != 1 {
- return nil, fmt.Errorf("Service length from meta db is %d", len(lists))
+ return nil, fmt.Errorf("service length from meta db is %d", len(lists))
}
var s v1.Service
@@ -154,3 +155,55 @@ func handleServiceFromMetaManager(content []byte) (*v1.Service, error) {
}
return &s, nil
}
+
+func (s *services) ListAll() ([]v1.Service, error) {
+ resource := fmt.Sprintf("%s/%s", s.namespace, constants.ResourceTypeService)
+ msg := message.BuildMsg(modules.MetaGroup, "", constant.ModuleNameEdgeMesh, resource, model.QueryOperation, nil)
+ respMsg, err := s.send.SendSync(msg)
+ if err != nil {
+ return nil, fmt.Errorf("get service list from metaManager failed, err: %v", err)
+ }
+ var content []byte
+ switch respMsg.Content.(type) {
+ case []byte:
+ content = respMsg.GetContent().([]byte)
+ default:
+ content, err = json.Marshal(respMsg.Content)
+ if err != nil {
+ return nil, fmt.Errorf("marshal message to service list failed, err: %v", err)
+ }
+ }
+
+ if respMsg.GetOperation() == model.ResponseOperation {
+ return handleServiceListFromMetaDB(content)
+ }
+ return handleServiceListFromMetaManager(content)
+}
+
+func handleServiceListFromMetaDB(content []byte) ([]v1.Service, error) {
+ var lists []string
+ err := json.Unmarshal([]byte(content), &lists)
+ if err != nil {
+ return nil, fmt.Errorf("unmarshal message to service list from edge db failed, err: %v", err)
+ }
+
+ var serviceList []v1.Service
+ for i := range lists {
+ var s v1.Service
+ err = json.Unmarshal([]byte(lists[i]), &s)
+ if err != nil {
+ return nil, fmt.Errorf("unmarshal message to service from edge db failed, err: %v", err)
+ }
+ serviceList = append(serviceList, s)
+ }
+ return serviceList, nil
+}
+
+func handleServiceListFromMetaManager(content []byte) ([]v1.Service, error) {
+ var serviceList []v1.Service
+ err := json.Unmarshal(content, &serviceList)
+ if err != nil {
+ return nil, fmt.Errorf("unmarshal message to service list failed, err: %v", err)
+ }
+ return serviceList, nil
+}
diff --git a/edge/pkg/metamanager/process.go b/edge/pkg/metamanager/process.go
index f527c03db..8cabaf359 100644
--- a/edge/pkg/metamanager/process.go
+++ b/edge/pkg/metamanager/process.go
@@ -97,6 +97,14 @@ func requireRemoteQuery(resType string) bool {
resType == model.ResourceTypeNode
}
+// if resource type is EdgeMesh related
+func isEdgeMeshResource(resType string) bool {
+ return resType == constants.ResourceTypeService ||
+ resType == constants.ResourceTypeServiceList ||
+ resType == constants.ResourceTypeEndpoints ||
+ resType == model.ResourceTypePodlist
+}
+
func isConnected() bool {
return metaManagerConfig.Connected
}
@@ -117,7 +125,6 @@ func resourceUnchanged(resType string, resKey string, content []byte) bool {
}
func (m *metaManager) processInsert(message model.Message) {
-
var err error
var content []byte
switch message.GetContent().(type) {
@@ -132,19 +139,53 @@ func (m *metaManager) processInsert(message model.Message) {
}
}
resKey, resType, _ := parseResource(message.GetResource())
+ switch resType {
+ case constants.ResourceTypeServiceList:
+ var svcList []v1.Service
+ err = json.Unmarshal(content, &svcList)
+ if err != nil {
+ klog.Errorf("Unmarshal insert message content failed, %s", msgDebugInfo(&message))
+ feedbackError(err, "Error to unmarshal", message)
+ return
+ }
+ for _, svc := range svcList {
+ data, err := json.Marshal(svc)
+ if err != nil {
+ klog.Errorf("Marshal service content failed, %v", svc)
+ continue
+ }
+ meta := &dao.Meta{
+ Key: fmt.Sprintf("%s/%s/%s", svc.Namespace, constants.ResourceTypeService, svc.Name),
+ Type: constants.ResourceTypeService,
+ Value: string(data)}
+ err = dao.SaveMeta(meta)
+ if err != nil {
+ klog.Errorf("Save meta %s failed, svc: %v, err: %v", string(data), svc, err)
+ feedbackError(err, "Error to save meta to DB", message)
+ return
+ }
+ }
+ default:
+ meta := &dao.Meta{
+ Key: resKey,
+ Type: resType,
+ Value: string(content)}
+ err = dao.SaveMeta(meta)
+ if err != nil {
+ klog.Errorf("save meta failed, %s: %v", msgDebugInfo(&message), err)
+ feedbackError(err, "Error to save meta to DB", message)
+ return
+ }
+ }
- meta := &dao.Meta{
- Key: resKey,
- Type: resType,
- Value: string(content)}
- err = dao.SaveMeta(meta)
- if err != nil {
- klog.Errorf("save meta failed, %s: %v", msgDebugInfo(&message), err)
- feedbackError(err, "Error to save meta to DB", message)
+ if resType == constants.ResourceTypeListener {
+ // Notify edgemesh only
+ resp := message.NewRespByMessage(&message, nil)
+ sendToEdgeMesh(resp, true)
return
}
- if resType == constants.ResourceTypeService || resType == constants.ResourceTypeEndpoints {
+ if isEdgeMeshResource(resType) {
// Notify edgemesh
sendToEdgeMesh(&message, false)
} else {
@@ -279,7 +320,7 @@ func (m *metaManager) processUpdate(message model.Message) {
resp := message.NewRespByMessage(&message, OK)
sendToEdged(resp, message.IsSync())
case CloudControlerModel:
- if resType == constants.ResourceTypeService || resType == constants.ResourceTypeEndpoints {
+ if isEdgeMeshResource(resType) {
sendToEdgeMesh(&message, message.IsSync())
} else {
sendToEdged(&message, message.IsSync())
@@ -296,7 +337,6 @@ func (m *metaManager) processUpdate(message model.Message) {
}
func (m *metaManager) processResponse(message model.Message) {
-
var err error
var content []byte
switch message.GetContent().(type) {
@@ -323,7 +363,7 @@ func (m *metaManager) processResponse(message model.Message) {
return
}
- // Notify edged or edgemesh if the data if coming from cloud
+ // Notify edged or edgemesh if the data is coming from cloud
if message.GetSource() == CloudControlerModel {
if resType == constants.ResourceTypeService || resType == constants.ResourceTypeEndpoints {
sendToEdgeMesh(&message, message.IsSync())
@@ -352,6 +392,12 @@ func (m *metaManager) processDelete(message model.Message) {
sendToCloud(resp)
return
}
+ if resType == constants.ResourceTypeListener {
+ // Notify edgemesh only
+ resp := message.NewRespByMessage(&message, OK)
+ sendToEdgeMesh(resp, true)
+ return
+ }
// Notify edged
sendToEdged(&message, false)
resp := message.NewRespByMessage(&message, OK)
@@ -386,7 +432,7 @@ func (m *metaManager) processQuery(message model.Message) {
} else {
resp := message.NewRespByMessage(&message, *metas)
resp.SetRoute(MetaManagerModuleName, resp.GetGroup())
- if resType == constants.ResourceTypeService || resType == constants.ResourceTypeEndpoints || resType == model.ResourceTypePodlist {
+ if resType == constants.ResourceTypeService || resType == constants.ResourceTypeEndpoints || resType == constants.ResourceTypeListener {
sendToEdgeMesh(resp, message.IsSync())
} else {
sendToEdged(resp, message.IsSync())