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, ©Data)
- 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(©Data)
- indexer = backend.Store().ServiceAlias()
- } else {
- prefix = apt.GenerateServiceIndexKey(©Data)
- 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
-}