diff options
| author | liuzhiyi1993 <liuzhiyi@huawei.com> | 2020-02-13 05:56:21 +0800 |
|---|---|---|
| committer | liuzhiyi1993 <liuzhiyi@huawei.com> | 2020-03-31 22:30:02 +0800 |
| commit | 0dfe63446ad76e41a084ade9d3c393a2554def34 (patch) | |
| tree | f9857c1fd4ad3e9a00aa04effdb1c57bde6398b7 /edge | |
| parent | Merge pull request #1567 from chendave/broken (diff) | |
| download | kubeedge-0dfe63446ad76e41a084ade9d3c393a2554def34.tar.gz | |
refactor edgemesh
Diffstat (limited to 'edge')
| -rw-r--r-- | edge/pkg/metamanager/client/listener.go | 78 | ||||
| -rw-r--r-- | edge/pkg/metamanager/client/metaclient.go | 6 | ||||
| -rw-r--r-- | edge/pkg/metamanager/client/service.go | 71 | ||||
| -rw-r--r-- | edge/pkg/metamanager/process.go | 74 |
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()) |
