This is an automated email from the ASF dual-hosted git repository.

tianxiaoliang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/servicecomb-service-center.git


The following commit(s) were added to refs/heads/master by this push:
     new 62509e4  [enhancement] optimize service/util func invoker (#708)
62509e4 is described below

commit 62509e43c39159a7eeaf8d5a19f9fbc566a23244
Author: popozy <[email protected]>
AuthorDate: Wed Oct 14 16:43:11 2020 +0800

    [enhancement] optimize service/util func invoker (#708)
    
    + keep util funcions under service/util still used by ms/etcd as 
serviceUtil.xxx
    + for different implementation of datasource interface, such as  mongo, 
util funcs can works under directory service/util/mongo
---
 server/service/dep/etcd/util.go |  20 ----
 server/service/ms/etcd/etcd.go  |  31 +++---
 server/service/ms/etcd/util.go  | 226 +---------------------------------------
 3 files changed, 15 insertions(+), 262 deletions(-)

diff --git a/server/service/dep/etcd/util.go b/server/service/dep/etcd/util.go
index 96b083c..e631e63 100644
--- a/server/service/dep/etcd/util.go
+++ b/server/service/dep/etcd/util.go
@@ -14,23 +14,3 @@
 // limitations under the License.
 
 package etcd
-
-import (
-       "encoding/json"
-       pb "github.com/apache/servicecomb-service-center/pkg/registry"
-       apt "github.com/apache/servicecomb-service-center/server/core"
-       "github.com/apache/servicecomb-service-center/server/plugin/registry"
-)
-
-func DeleteDependencyForDeleteService(domainProject string, serviceID string, 
service *pb.MicroServiceKey) (registry.PluginOp, error) {
-       key := apt.GenerateConsumerDependencyQueueKey(domainProject, serviceID, 
apt.DepsQueueUUID)
-       conDep := new(pb.ConsumerDependency)
-       conDep.Consumer = service
-       conDep.Providers = []*pb.MicroServiceKey{}
-       conDep.Override = true
-       data, err := json.Marshal(conDep)
-       if err != nil {
-               return registry.PluginOp{}, err
-       }
-       return registry.OpPut(registry.WithStrKey(key), 
registry.WithValue(data)), nil
-}
diff --git a/server/service/ms/etcd/etcd.go b/server/service/ms/etcd/etcd.go
index a9671f8..aba1669 100644
--- a/server/service/ms/etcd/etcd.go
+++ b/server/service/ms/etcd/etcd.go
@@ -30,7 +30,6 @@ import (
        "github.com/apache/servicecomb-service-center/server/plugin/quota"
        "github.com/apache/servicecomb-service-center/server/plugin/registry"
        scerr "github.com/apache/servicecomb-service-center/server/scerror"
-       depUtil 
"github.com/apache/servicecomb-service-center/server/service/dep/etcd"
        "github.com/apache/servicecomb-service-center/server/service/ms"
        serviceUtil 
"github.com/apache/servicecomb-service-center/server/service/util"
        "strconv"
@@ -92,9 +91,7 @@ func (ds *DataSource) RegisterService(ctx context.Context, 
request *pb.CreateSer
 
 func (ds *DataSource) GetServices(ctx context.Context, request 
*pb.GetServicesRequest) (
        *pb.GetServicesResponse, error) {
-       domainProject := util.ParseDomainProject(ctx)
-
-       services, err := getAllServiceUtil(ctx, domainProject)
+       services, err := serviceUtil.GetAllServiceUtil(ctx)
        if err != nil {
                log.Errorf(err, "get all services by domain failed")
                return &pb.GetServicesResponse{
@@ -111,7 +108,7 @@ func (ds *DataSource) GetServices(ctx context.Context, 
request *pb.GetServicesRe
 func (ds *DataSource) GetService(ctx context.Context, request 
*pb.GetServiceRequest) (
        *pb.GetServiceResponse, error) {
        domainProject := util.ParseDomainProject(ctx)
-       singleService, err := getSingleService(ctx, domainProject, 
request.ServiceId)
+       singleService, err := serviceUtil.GetService(ctx, domainProject, 
request.ServiceId)
 
        if err != nil {
                log.Errorf(err, "get micro-service[%s] failed, get service file 
failed", request.ServiceId)
@@ -137,7 +134,7 @@ func (ds *DataSource) ExistService(ctx context.Context, 
request *pb.GetExistence
        serviceFlag := util.StringJoin([]string{
                request.Environment, request.AppId, request.ServiceName, 
request.Version}, "/")
 
-       ids, exist, err := findServiceIds(ctx, request.Version, 
&pb.MicroServiceKey{
+       ids, exist, err := serviceUtil.FindServiceIds(ctx, request.Version, 
&pb.MicroServiceKey{
                Environment: request.Environment,
                AppId:       request.AppId,
                ServiceName: request.ServiceName,
@@ -175,7 +172,7 @@ func (ds *DataSource) UpdateService(ctx context.Context, 
request *pb.UpdateServi
        domainProject := util.ParseDomainProject(ctx)
 
        key := apt.GenerateServiceKey(domainProject, request.ServiceId)
-       microservice, err := getService(ctx, domainProject, request.ServiceId)
+       microservice, err := serviceUtil.GetService(ctx, domainProject, 
request.ServiceId)
        if err != nil {
                log.Errorf(err, "update service[%s] properties failed, get 
service file failed, operator: %s",
                        request.ServiceId, remoteIP)
@@ -401,7 +398,7 @@ func (ds *DataSource) Heartbeat(ctx context.Context, 
request *pb.HeartbeatReques
        domainProject := util.ParseDomainProject(ctx)
        instanceFlag := util.StringJoin([]string{request.ServiceId, 
request.InstanceId}, "/")
 
-       _, ttl, err := HeartbeatUtil(ctx, domainProject, request.ServiceId, 
request.InstanceId)
+       _, ttl, err := serviceUtil.HeartbeatUtil(ctx, domainProject, 
request.ServiceId, request.InstanceId)
        if err != nil {
                log.Errorf(err, "heartbeat failed, instance[%s]. operator %s",
                        instanceFlag, remoteIP)
@@ -431,7 +428,7 @@ func (ds *DataSource) ModifySchemas(ctx context.Context, 
request *pb.ModifySchem
        serviceID := request.ServiceId
        domainProject := util.ParseDomainProject(ctx)
 
-       serviceInfo, err := getService(ctx, domainProject, serviceID)
+       serviceInfo, err := serviceUtil.GetService(ctx, domainProject, 
serviceID)
        if err != nil {
                log.Errorf(err, "modify service[%s] schemas failed, get service 
failed, operator: %s", serviceID, remoteIP)
                return &pb.ModifySchemasResponse{
@@ -504,7 +501,7 @@ func (ds *DataSource) ExistSchema(ctx context.Context, 
request *pb.GetExistenceR
        *pb.GetExistenceResponse, error) {
        domainProject := util.ParseDomainProject(ctx)
 
-       if !serviceExist(ctx, domainProject, request.ServiceId) {
+       if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
                log.Warnf("schema[%s/%s] exist failed, service does not exist", 
request.ServiceId, request.SchemaId)
                return &pb.GetExistenceResponse{
                        Response: 
proto.CreateResponse(scerr.ErrServiceNotExists, "service does not exist."),
@@ -598,7 +595,7 @@ func (ds *DataSource) modifySchemas(ctx context.Context, 
domainProject string, s
                        }
 
                        service.Schemas = nonExistSchemaIds
-                       opt, err := updateService(domainProject, serviceID, 
service)
+                       opt, err := serviceUtil.UpdateService(domainProject, 
serviceID, service)
                        if err != nil {
                                log.Errorf(err, "modify service[%s] schemas 
failed, update service.Schemas failed, operator: %s",
                                        serviceID, remoteIP)
@@ -665,7 +662,7 @@ func (ds *DataSource) modifySchemas(ctx context.Context, 
domainProject string, s
                }
 
                service.Schemas = schemaIDs
-               opt, err := updateService(domainProject, serviceID, service)
+               opt, err := serviceUtil.UpdateService(domainProject, serviceID, 
service)
                if err != nil {
                        log.Errorf(err, "modify service[%s] schemas failed, 
update service.Schemas failed, operator: %s",
                                serviceID, remoteIP)
@@ -699,7 +696,7 @@ func (ds *DataSource) modifySchema(ctx context.Context, 
serviceID string, schema
        domainProject := util.ParseDomainProject(ctx)
        schemaID := schema.SchemaId
 
-       microService, err := getService(ctx, domainProject, serviceID)
+       microService, err := serviceUtil.GetService(ctx, domainProject, 
serviceID)
        if err != nil {
                log.Errorf(err, "modify schema[%s/%s] failed, get `microService 
failed, operator: %s",
                        serviceID, schemaID, remoteIP)
@@ -750,7 +747,7 @@ func (ds *DataSource) modifySchema(ctx context.Context, 
serviceID string, schema
 
                if len(microService.Schemas) == 0 {
                        microService.Schemas = append(microService.Schemas, 
schemaID)
-                       opt, err := updateService(domainProject, serviceID, 
microService)
+                       opt, err := serviceUtil.UpdateService(domainProject, 
serviceID, microService)
                        if err != nil {
                                log.Errorf(err, "modify schema[%s/%s] failed, 
update microService.Schemas failed, operator: %s",
                                        serviceID, schemaID, remoteIP)
@@ -761,7 +758,7 @@ func (ds *DataSource) modifySchema(ctx context.Context, 
serviceID string, schema
        } else {
                if !isExist {
                        microService.Schemas = append(microService.Schemas, 
schemaID)
-                       opt, err := updateService(domainProject, serviceID, 
microService)
+                       opt, err := serviceUtil.UpdateService(domainProject, 
serviceID, microService)
                        if err != nil {
                                log.Errorf(err, "modify schema[%s/%s] failed, 
update microService.Schemas failed, operator: %s",
                                        serviceID, schemaID, remoteIP)
@@ -803,7 +800,7 @@ func (ds *DataSource) DeleteServicePri(ctx context.Context, 
serviceID string, fo
                return proto.CreateResponse(scerr.ErrInvalidParams, 
err.Error()), nil
        }
 
-       microservice, err := getService(ctx, domainProject, serviceID)
+       microservice, err := serviceUtil.GetService(ctx, domainProject, 
serviceID)
        if err != nil {
                log.Errorf(err, "%s micro-service[%s] failed, get service file 
failed, operator: %s",
                        title, serviceID, remoteIP)
@@ -865,7 +862,7 @@ func (ds *DataSource) DeleteServicePri(ctx context.Context, 
serviceID string, fo
        }
 
        //删除依赖规则
-       optDeleteDep, err := 
depUtil.DeleteDependencyForDeleteService(domainProject, serviceID, serviceKey)
+       optDeleteDep, err := 
serviceUtil.DeleteDependencyForDeleteService(domainProject, serviceID, 
serviceKey)
        if err != nil {
                log.Errorf(err, "%s micro-service[%s] failed, delete dependency 
failed, operator: %s",
                        title, serviceID, remoteIP)
diff --git a/server/service/ms/etcd/util.go b/server/service/ms/etcd/util.go
index ef8e321..83736b7 100644
--- a/server/service/ms/etcd/util.go
+++ b/server/service/ms/etcd/util.go
@@ -18,7 +18,6 @@ package etcd
 import (
        "context"
        "encoding/json"
-       "errors"
        "github.com/apache/servicecomb-service-center/pkg/log"
        pb "github.com/apache/servicecomb-service-center/pkg/registry"
        "github.com/apache/servicecomb-service-center/pkg/util"
@@ -26,7 +25,6 @@ import (
        "github.com/apache/servicecomb-service-center/server/core/backend"
        "github.com/apache/servicecomb-service-center/server/core/proto"
        "github.com/apache/servicecomb-service-center/server/plugin"
-       "github.com/apache/servicecomb-service-center/server/plugin/discovery"
        "github.com/apache/servicecomb-service-center/server/plugin/registry"
        "github.com/apache/servicecomb-service-center/server/plugin/uuid"
        scerr "github.com/apache/servicecomb-service-center/server/scerror"
@@ -36,56 +34,7 @@ import (
        "time"
 )
 
-var (
-       ErrLeaseIDNotExist = errors.New("leaseId not exist, instance not exist")
-)
-
 // service
-func getSingleService(ctx context.Context, domainProject string, serviceID 
string) (*pb.MicroService, error) {
-       key := apt.GenerateServiceKey(domainProject, serviceID)
-       opts := append(serviceUtil.FromContext(ctx), registry.WithStrKey(key))
-       serviceResp, err := backend.Store().Service().Search(ctx, opts...)
-       if err != nil {
-               return nil, err
-       }
-       if len(serviceResp.Kvs) == 0 {
-               return nil, nil
-       }
-       return serviceResp.Kvs[0].Value.(*pb.MicroService), nil
-}
-
-func getAllServiceUtil(ctx context.Context, domainProject string) 
([]*pb.MicroService, error) {
-       services, err := getServicesByDomainProject(ctx, domainProject)
-       if err != nil {
-               return nil, err
-       }
-       return services, nil
-}
-
-func getServicesByDomainProject(ctx context.Context, domainProject string) 
([]*pb.MicroService, error) {
-       kvs, err := getServicesRawData(ctx, domainProject)
-       if err != nil {
-               return nil, err
-       }
-       var services []*pb.MicroService
-       for _, kv := range kvs {
-               services = append(services, kv.Value.(*pb.MicroService))
-       }
-       return services, nil
-}
-
-func getServicesRawData(ctx context.Context, domainProject string) 
([]*discovery.KeyValue, error) {
-       key := apt.GenerateServiceKey(domainProject, "")
-       opts := append(serviceUtil.FromContext(ctx),
-               registry.WithStrKey(key),
-               registry.WithPrefix())
-       resp, err := backend.Store().Service().Search(ctx, opts...)
-       if err != nil {
-               return nil, err
-       }
-       return resp.Kvs, err
-}
-
 func capRegisterData(ctx context.Context, request *pb.CreateServiceRequest) (
        []registry.PluginOp, []registry.CompareOp, []registry.PluginOp, error) {
        remoteIP := util.GetIPFromContext(ctx)
@@ -200,125 +149,6 @@ func newRegisterServiceResp(ctx context.Context, 
reqService *pb.MicroService, re
        }, nil
 }
 
-func serviceExist(ctx context.Context, domainProject string, serviceID string) 
bool {
-       opts := append(serviceUtil.FromContext(ctx),
-               registry.WithStrKey(apt.GenerateServiceKey(domainProject, 
serviceID)),
-               registry.WithCountOnly())
-       resp, err := backend.Store().Service().Search(ctx, opts...)
-       if err != nil || resp.Count == 0 {
-               return false
-       }
-       return true
-}
-
-func findServiceIds(ctx context.Context, versionRule string, key 
*pb.MicroServiceKey) ([]string, bool, error) {
-       // 版本规则
-       match := serviceUtil.ParseVersionRule(versionRule)
-       if match == nil {
-               copyData := *key
-               copyData.Version = versionRule
-               serviceID, err := getServiceID(ctx, &copyData)
-               if err != nil {
-                       return nil, false, err
-               }
-               if len(serviceID) > 0 {
-                       return []string{serviceID}, true, nil
-               }
-               return nil, false, nil
-       }
-
-       searchAlias := false
-       alsoFindAlias := len(key.Alias) > 0
-
-FindRule:
-       resp, err := getServiceAllVersions(ctx, key, searchAlias)
-       if err != nil {
-               return nil, false, err
-       }
-       if len(resp.Kvs) == 0 {
-               if !alsoFindAlias {
-                       return nil, false, nil
-               }
-               searchAlias = true
-               alsoFindAlias = false
-               goto FindRule
-       }
-       return match(resp.Kvs), true, nil
-}
-
-func getServiceID(ctx context.Context, key *pb.MicroServiceKey) (serviceID 
string, err error) {
-       serviceID, err = searchServiceID(ctx, key)
-       if err != nil {
-               return
-       }
-       if len(serviceID) == 0 {
-               // 别名查询
-               log.Debugf("could not search microservice[%s/%s/%s/%s] id by 
'serviceName', now try 'alias'",
-                       key.Environment, key.AppId, key.ServiceName, 
key.Version)
-               return searchServiceIDFromAlias(ctx, key)
-       }
-       return
-}
-
-func searchServiceID(ctx context.Context, key *pb.MicroServiceKey) (string, 
error) {
-       opts := append(serviceUtil.FromContext(ctx), 
registry.WithStrKey(apt.GenerateServiceIndexKey(key)))
-       resp, err := backend.Store().ServiceIndex().Search(ctx, opts...)
-       if err != nil {
-               return "", err
-       }
-       if len(resp.Kvs) == 0 {
-               return "", nil
-       }
-       return resp.Kvs[0].Value.(string), nil
-}
-
-func searchServiceIDFromAlias(ctx context.Context, key *pb.MicroServiceKey) 
(string, error) {
-       opts := append(serviceUtil.FromContext(ctx), 
registry.WithStrKey(apt.GenerateServiceAliasKey(key)))
-       resp, err := backend.Store().ServiceAlias().Search(ctx, opts...)
-       if err != nil {
-               return "", err
-       }
-       if len(resp.Kvs) == 0 {
-               return "", nil
-       }
-       return resp.Kvs[0].Value.(string), nil
-}
-
-func getServiceAllVersions(ctx context.Context, key *pb.MicroServiceKey, alias 
bool) (*discovery.Response, error) {
-       copyData := *key
-       copyData.Version = ""
-       var (
-               prefix  string
-               indexer discovery.Indexer
-       )
-       if alias {
-               prefix = apt.GenerateServiceAliasKey(&copyData)
-               indexer = backend.Store().ServiceAlias()
-       } else {
-               prefix = apt.GenerateServiceIndexKey(&copyData)
-               indexer = backend.Store().ServiceIndex()
-       }
-       opts := append(serviceUtil.FromContext(ctx),
-               registry.WithStrKey(prefix),
-               registry.WithPrefix(),
-               registry.WithDescendOrder())
-       resp, err := indexer.Search(ctx, opts...)
-       return resp, err
-}
-
-func updateService(domainProject string, serviceID string, service 
*pb.MicroService) (opt registry.PluginOp,
-       err error) {
-       opt = registry.PluginOp{}
-       key := apt.GenerateServiceKey(domainProject, serviceID)
-       data, err := json.Marshal(service)
-       if err != nil {
-               log.Errorf(err, "marshal service file failed")
-               return
-       }
-       opt = registry.OpPut(registry.WithStrKey(key), registry.WithValue(data))
-       return
-}
-
 // schema
 func getSchemaSummary(ctx context.Context, domainProject string, serviceID 
string, schemaID string) (string, error) {
        key := apt.GenerateServiceSchemaSummaryKey(domainProject, serviceID, 
schemaID)
@@ -335,19 +165,6 @@ func getSchemaSummary(ctx context.Context, domainProject 
string, serviceID strin
        return resp.Kvs[0].Value.(string), nil
 }
 
-func getService(ctx context.Context, domainProject string, serviceID string) 
(*pb.MicroService, error) {
-       key := apt.GenerateServiceKey(domainProject, serviceID)
-       opts := append(serviceUtil.FromContext(ctx), registry.WithStrKey(key))
-       serviceResp, err := backend.Store().Service().Search(ctx, opts...)
-       if err != nil {
-               return nil, err
-       }
-       if len(serviceResp.Kvs) == 0 {
-               return nil, nil
-       }
-       return serviceResp.Kvs[0].Value.(*pb.MicroService), nil
-}
-
 func getSchemasFromDatabase(ctx context.Context, domainProject string, 
serviceID string) ([]*pb.Schema, error) {
        key := apt.GenerateServiceSchemaKey(domainProject, serviceID, "")
        resp, err := backend.Store().Schema().Search(ctx,
@@ -493,18 +310,6 @@ func commitSchemaInfo(domainProject string, serviceID 
string, schema *pb.Schema)
 }
 
 // instance util
-func HeartbeatUtil(ctx context.Context, domainProject string, serviceID 
string, instanceID string) (leaseID int64, ttl int64, _ *scerr.Error) {
-       leaseID, err := GetLeaseID(ctx, domainProject, serviceID, instanceID)
-       if err != nil {
-               return leaseID, ttl, 
scerr.NewError(scerr.ErrUnavailableBackend, err.Error())
-       }
-       ttl, err = KeepAliveLease(ctx, domainProject, serviceID, instanceID, 
leaseID)
-       if err != nil {
-               return leaseID, ttl, scerr.NewError(scerr.ErrInstanceNotExists, 
err.Error())
-       }
-       return leaseID, ttl, nil
-}
-
 func preProcessRegisterInstance(ctx context.Context, instance 
*pb.MicroServiceInstance) *scerr.Error {
        if len(instance.Status) == 0 {
                instance.Status = pb.MSI_UP
@@ -542,39 +347,10 @@ func preProcessRegisterInstance(ctx context.Context, 
instance *pb.MicroServiceIn
        }
 
        domainProject := util.ParseDomainProject(ctx)
-       microservice, err := getService(ctx, domainProject, instance.ServiceId)
+       microservice, err := serviceUtil.GetService(ctx, domainProject, 
instance.ServiceId)
        if microservice == nil || err != nil {
                return scerr.NewError(scerr.ErrServiceNotExists, "Invalid 
'serviceID' in request body.")
        }
        instance.Version = microservice.Version
        return nil
 }
-
-// heartbeat util
-func GetLeaseID(ctx context.Context, domainProject string, serviceID string, 
instanceID string) (int64, error) {
-       opts := append(serviceUtil.FromContext(ctx),
-               registry.WithStrKey(apt.GenerateInstanceLeaseKey(domainProject, 
serviceID, instanceID)))
-       resp, err := backend.Store().Lease().Search(ctx, opts...)
-       if err != nil {
-               return -1, err
-       }
-       if len(resp.Kvs) <= 0 {
-               return -1, nil
-       }
-       leaseID, _ := strconv.ParseInt(resp.Kvs[0].Value.(string), 10, 64)
-       return leaseID, nil
-}
-
-func KeepAliveLease(ctx context.Context, domainProject, serviceID, instanceID 
string, leaseID int64) (
-       ttl int64, err error) {
-       if leaseID == -1 {
-               return ttl, ErrLeaseIDNotExist
-       }
-       ttl, err = backend.Store().KeepAlive(ctx,
-               registry.WithStrKey(apt.GenerateInstanceLeaseKey(domainProject, 
serviceID, instanceID)),
-               registry.WithLease(leaseID))
-       if err != nil {
-               return ttl, err
-       }
-       return ttl, nil
-}

Reply via email to