This is an automated email from the ASF dual-hosted git repository.
littlecui 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 baf44a7 [SCB-2094] Implement server/rest/govern exist interface (#725)
baf44a7 is described below
commit baf44a75e4ca4cfe8101b823bea411f1b529762c
Author: xzccfzy <[email protected]>
AuthorDate: Wed Oct 21 10:33:08 2020 +0800
[SCB-2094] Implement server/rest/govern exist interface (#725)
Co-authored-by: 薛泽超 <[email protected]>
---
datasource/etcd/dep.go | 12 +-
datasource/etcd/key_generator.go | 82 +++++++++++-
datasource/etcd/ms.go | 182 +++++++++++++++++++++++++++
datasource/etcd/ms_test.go | 151 ++++++++++++++++++++++
datasource/etcd/util.go | 263 +++++++++++++++++++++++++++++++++++++++
datasource/ms.go | 5 +
6 files changed, 679 insertions(+), 16 deletions(-)
diff --git a/datasource/etcd/dep.go b/datasource/etcd/dep.go
index aad3796..5e04fd4 100644
--- a/datasource/etcd/dep.go
+++ b/datasource/etcd/dep.go
@@ -113,7 +113,7 @@ func (ds *DataSource) AddOrUpdateDependencies(ctx
context.Context, dependencyInf
if !override {
id = util.GenerateUUID()
}
- key := apt.GenerateConsumerDependencyQueueKey(domainProject,
consumerID, id)
+ key := GenerateConsumerDependencyQueueKey(domainProject,
consumerID, id)
opts = append(opts, registry.OpPut(registry.WithStrKey(key),
registry.WithValue(data)))
}
@@ -127,13 +127,3 @@ func (ds *DataSource) AddOrUpdateDependencies(ctx
context.Context, dependencyInf
override, dependencyInfos, util.GetIPFromContext(ctx))
return proto.CreateResponse(proto.Response_SUCCESS, "Create dependency
successfully."), nil
}
-
-func toDependencyFilterOptions(in *pb.GetDependenciesRequest) (opts
[]serviceUtil.DependencyRelationFilterOption) {
- if in.SameDomain {
- opts = append(opts, serviceUtil.WithSameDomainProject())
- }
- if in.NoSelf {
- opts = append(opts, serviceUtil.WithoutSelfDependency())
- }
- return opts
-}
diff --git a/datasource/etcd/key_generator.go b/datasource/etcd/key_generator.go
index 486dba8..9e7df66 100644
--- a/datasource/etcd/key_generator.go
+++ b/datasource/etcd/key_generator.go
@@ -1,12 +1,20 @@
package etcd
-import "github.com/apache/servicecomb-service-center/pkg/util"
+import (
+ "fmt"
+ "github.com/apache/servicecomb-service-center/pkg/registry"
+ "github.com/apache/servicecomb-service-center/pkg/util"
+ "strings"
+)
const (
- RegistryRootKey = "cse-sr"
- RegistryProjectKey = "projects"
- SPLIT = "/"
- RegistryDomainKey = "domains"
+ RegistryRootKey = "cse-sr"
+ RegistryProjectKey = "projects"
+ SPLIT = "/"
+ RegistryDomainKey = "domains"
+ RegistryServiceKey = "ms"
+ RegistryIndex = "indexes"
+ RegistryDepsQueueKey = "dep-queue"
)
func GetRootKey() string {
@@ -57,3 +65,67 @@ func GetDomainRootKey() string {
RegistryDomainKey,
}, SPLIT)
}
+
+func GenerateServiceIndexKey(key *registry.MicroServiceKey) string {
+ return util.StringJoin([]string{
+ GetServiceIndexRootKey(key.Tenant),
+ key.Environment,
+ key.AppId,
+ key.ServiceName,
+ key.Version,
+ }, SPLIT)
+}
+
+func GetServiceIndexRootKey(domainProject string) string {
+ return util.StringJoin([]string{
+ GetRootKey(),
+ RegistryServiceKey,
+ RegistryIndex,
+ domainProject,
+ }, SPLIT)
+}
+
+func KvToResponse(key []byte) (keys []string) {
+ return strings.Split(util.BytesToStringWithNoCopy(key), SPLIT)
+}
+
+func GetInfoFromSvcIndexKV(key []byte) *registry.MicroServiceKey {
+ keys := KvToResponse(key)
+ l := len(keys)
+ if l < 6 {
+ return nil
+ }
+ domainProject := fmt.Sprintf("%s/%s", keys[l-6], keys[l-5])
+ return ®istry.MicroServiceKey{
+ Tenant: domainProject,
+ Environment: keys[l-4],
+ AppId: keys[l-3],
+ ServiceName: keys[l-2],
+ Version: keys[l-1],
+ }
+}
+
+func GetServiceAppKey(domainProject, env, appID string) string {
+ return util.StringJoin([]string{
+ GetServiceIndexRootKey(domainProject),
+ env,
+ appID,
+ }, SPLIT)
+}
+
+func GenerateConsumerDependencyQueueKey(domainProject, consumerID, uuid
string) string {
+ return util.StringJoin([]string{
+ GetServiceDependencyQueueRootKey(domainProject),
+ consumerID,
+ uuid,
+ }, SPLIT)
+}
+
+func GetServiceDependencyQueueRootKey(domainProject string) string {
+ return util.StringJoin([]string{
+ GetRootKey(),
+ RegistryServiceKey,
+ RegistryDepsQueueKey,
+ domainProject,
+ }, SPLIT)
+}
diff --git a/datasource/etcd/ms.go b/datasource/etcd/ms.go
index cabf9f1..d7793e8 100644
--- a/datasource/etcd/ms.go
+++ b/datasource/etcd/ms.go
@@ -100,6 +100,188 @@ func (ds *DataSource) GetService(ctx context.Context,
request *pb.GetServiceRequ
}, nil
}
+func (ds *DataSource) GetServiceDetail(ctx context.Context, request
*pb.GetServiceRequest) (
+ *pb.GetServiceDetailResponse, error) {
+ ctx = util.SetContext(ctx, util.CtxCacheOnly, "1")
+
+ domainProject := util.ParseDomainProject(ctx)
+ options := []string{"tags", "rules", "instances", "schemas",
"dependencies"}
+
+ if len(request.ServiceId) == 0 {
+ return &pb.GetServiceDetailResponse{
+ Response: proto.CreateResponse(scerr.ErrInvalidParams,
"Invalid request for getting service detail."),
+ }, nil
+ }
+
+ service, err := serviceUtil.GetService(ctx, domainProject,
request.ServiceId)
+ if service == nil {
+ return &pb.GetServiceDetailResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+ if err != nil {
+ return &pb.GetServiceDetailResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ key := &pb.MicroServiceKey{
+ Tenant: domainProject,
+ Environment: service.Environment,
+ AppId: service.AppId,
+ ServiceName: service.ServiceName,
+ Version: "",
+ }
+ versions, err := getServiceAllVersions(ctx, key)
+ if err != nil {
+ log.Errorf(err, "get service[%s/%s/%s] all versions failed",
+ service.Environment, service.AppId, service.ServiceName)
+ return &pb.GetServiceDetailResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ serviceInfo, err := getServiceDetailUtil(ctx, ServiceDetailOpt{
+ domainProject: domainProject,
+ service: service,
+ options: options,
+ })
+ if err != nil {
+ return &pb.GetServiceDetailResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ serviceInfo.MicroService = service
+ serviceInfo.MicroServiceVersions = versions
+ return &pb.GetServiceDetailResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Get
service successfully."),
+ Service: serviceInfo,
+ }, nil
+}
+
+func (ds *DataSource) GetServicesInfo(ctx context.Context, request
*pb.GetServicesInfoRequest) (
+ *pb.GetServicesInfoResponse, error) {
+ ctx = util.SetContext(ctx, util.CtxCacheOnly, "1")
+
+ optionMap := make(map[string]struct{}, len(request.Options))
+ for _, opt := range request.Options {
+ optionMap[opt] = struct{}{}
+ }
+
+ options := make([]string, 0, len(optionMap))
+ if _, ok := optionMap["all"]; ok {
+ optionMap["statistics"] = struct{}{}
+ options = []string{"tags", "rules", "instances", "schemas",
"dependencies"}
+ } else {
+ for opt := range optionMap {
+ options = append(options, opt)
+ }
+ }
+
+ var st *pb.Statistics
+ if _, ok := optionMap["statistics"]; ok {
+ var err error
+ st, err = statistics(ctx, request.WithShared)
+ if err != nil {
+ return &pb.GetServicesInfoResponse{
+ Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
+ }, err
+ }
+ if len(optionMap) == 1 {
+ return &pb.GetServicesInfoResponse{
+ Response:
proto.CreateResponse(proto.Response_SUCCESS, "Statistics successfully."),
+ Statistics: st,
+ }, nil
+ }
+ }
+
+ //获取所有服务
+ services, err := serviceUtil.GetAllServiceUtil(ctx)
+ if err != nil {
+ log.Errorf(err, "get all services by domain failed")
+ return &pb.GetServicesInfoResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ allServiceDetails := make([]*pb.ServiceDetail, 0, len(services))
+ domainProject := util.ParseDomainProject(ctx)
+ for _, service := range services {
+ if !request.WithShared &&
apt.IsShared(proto.MicroServiceToKey(domainProject, service)) {
+ continue
+ }
+ if len(request.AppId) > 0 {
+ if request.AppId != service.AppId {
+ continue
+ }
+ if len(request.ServiceName) > 0 && request.ServiceName
!= service.ServiceName {
+ continue
+ }
+ }
+
+ serviceDetail, err := getServiceDetailUtil(ctx,
ServiceDetailOpt{
+ domainProject: domainProject,
+ service: service,
+ countOnly: request.CountOnly,
+ options: options,
+ })
+ if err != nil {
+ return &pb.GetServicesInfoResponse{
+ Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
+ }, err
+ }
+ serviceDetail.MicroService = service
+ allServiceDetails = append(allServiceDetails, serviceDetail)
+ }
+
+ return &pb.GetServicesInfoResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
"Get services info successfully."),
+ AllServicesDetail: allServiceDetails,
+ Statistics: st,
+ }, nil
+}
+
+func (ds *DataSource) GetApplications(ctx context.Context, request
*pb.GetAppsRequest) (*pb.GetAppsResponse, error) {
+ domainProject := util.ParseDomainProject(ctx)
+ key := GetServiceAppKey(domainProject, request.Environment, "")
+
+ opts := append(serviceUtil.FromContext(ctx),
+ registry.WithStrKey(key),
+ registry.WithPrefix(),
+ registry.WithKeyOnly())
+
+ resp, err := backend.Store().ServiceIndex().Search(ctx, opts...)
+ if err != nil {
+ return nil, err
+ }
+ l := len(resp.Kvs)
+ if l == 0 {
+ return &pb.GetAppsResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
"Get all applications successfully."),
+ }, nil
+ }
+
+ apps := make([]string, 0, l)
+ appMap := make(map[string]struct{}, l)
+ for _, kv := range resp.Kvs {
+ key := GetInfoFromSvcIndexKV(kv.Key)
+ if !request.WithShared && apt.IsShared(key) {
+ continue
+ }
+ if _, ok := appMap[key.AppId]; ok {
+ continue
+ }
+ appMap[key.AppId] = struct{}{}
+ apps = append(apps, key.AppId)
+ }
+
+ return &pb.GetAppsResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Get all
applications successfully."),
+ AppIds: apps,
+ }, nil
+}
+
func (ds *DataSource) ExistService(ctx context.Context, request
*pb.GetExistenceRequest) (*pb.GetExistenceResponse,
error) {
domainProject := util.ParseDomainProject(ctx)
diff --git a/datasource/etcd/ms_test.go b/datasource/etcd/ms_test.go
index 49e5d83..97b9886 100644
--- a/datasource/etcd/ms_test.go
+++ b/datasource/etcd/ms_test.go
@@ -737,6 +737,157 @@ func TestService_Delete(t *testing.T) {
})
}
+func TestService_Info(t *testing.T) {
+ datasource.Install("etcd", func(opts datasource.Options)
(datasource.DataSource, error) {
+ return NewDataSource(opts), nil
+ })
+
+ err := datasource.Init(datasource.Options{
+ Endpoint: "",
+ PluginImplName:
datasource.ImplName(archaius.GetString("servicecomb.datasource.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ t.Run("get all services", func(t *testing.T) {
+ log.Info("should be passed")
+ resp, err :=
datasource.Instance().GetServicesInfo(getContext(), &pb.GetServicesInfoRequest{
+ Options: []string{"all"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ resp, err = datasource.Instance().GetServicesInfo(getContext(),
&pb.GetServicesInfoRequest{
+ Options: []string{""},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ resp, err = datasource.Instance().GetServicesInfo(getContext(),
&pb.GetServicesInfoRequest{
+ Options: []string{"tags", "rules", "instances",
"schemas", "statistics"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ resp, err = datasource.Instance().GetServicesInfo(getContext(),
&pb.GetServicesInfoRequest{
+ Options: []string{"statistics"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ resp, err = datasource.Instance().GetServicesInfo(getContext(),
&pb.GetServicesInfoRequest{
+ Options: []string{"instances"},
+ CountOnly: true,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+}
+
+func TestService_Detail(t *testing.T) {
+ datasource.Install("etcd", func(opts datasource.Options)
(datasource.DataSource, error) {
+ return NewDataSource(opts), nil
+ })
+
+ err := datasource.Init(datasource.Options{
+ Endpoint: "",
+ PluginImplName:
datasource.ImplName(archaius.GetString("servicecomb.datasource.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ var (
+ serviceId string
+ )
+
+ t.Run("execute 'get detail' operation", func(t *testing.T) {
+ log.Info("should be passed")
+ resp, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "govern_service_group",
+ ServiceName: "govern_service_name",
+ Version: "3.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ serviceId = resp.ServiceId
+
+ datasource.Instance().ModifySchema(getContext(),
&pb.ModifySchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "schemaId",
+ Schema: "detail",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ datasource.Instance().RegisterInstance(getContext(),
&pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId,
+ Endpoints: []string{
+ "govern:127.0.0.1:8080",
+ },
+ HostName: "UT-HOST",
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ log.Info("when get invalid service detail, should be failed")
+ respD, err :=
datasource.Instance().GetServiceDetail(getContext(), &pb.GetServiceRequest{
+ ServiceId: "",
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respD.Response.GetCode())
+
+ log.Info("when get a service detail, should be passed")
+ respGetServiceDetail, err :=
datasource.Instance().GetServiceDetail(getContext(), &pb.GetServiceRequest{
+ ServiceId: serviceId,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respGetServiceDetail.Response.GetCode())
+
+ respDelete, err :=
datasource.Instance().UnregisterService(getContext(), &pb.DeleteServiceRequest{
+ ServiceId: serviceId,
+ Force: true,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respDelete.Response.GetCode())
+
+ respGetServiceDetail, err =
datasource.Instance().GetServiceDetail(getContext(), &pb.GetServiceRequest{
+ ServiceId: serviceId,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respGetServiceDetail.Response.GetCode())
+ })
+}
+
+func TestApplication_Get(t *testing.T) {
+ datasource.Install("etcd", func(opts datasource.Options)
(datasource.DataSource, error) {
+ return NewDataSource(opts), nil
+ })
+
+ err := datasource.Init(datasource.Options{
+ Endpoint: "",
+ PluginImplName:
datasource.ImplName(archaius.GetString("servicecomb.datasource.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ t.Run("execute 'get apps' operation", func(t *testing.T) {
+ log.Info("when request is valid, should be passed")
+ resp, err :=
datasource.Instance().GetApplications(getContext(), &pb.GetAppsRequest{})
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ resp, err = datasource.Instance().GetApplications(getContext(),
&pb.GetAppsRequest{
+ Environment: pb.ENV_ACCEPT,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+}
+
func TestInstance_Create(t *testing.T) {
datasource.Install("etcd", func(opts datasource.Options)
(datasource.DataSource, error) {
return NewDataSource(opts), nil
diff --git a/datasource/etcd/util.go b/datasource/etcd/util.go
index 320e52d..a9f1c3e 100644
--- a/datasource/etcd/util.go
+++ b/datasource/etcd/util.go
@@ -21,6 +21,7 @@ import (
"context"
"encoding/json"
errorsEx "github.com/apache/servicecomb-service-center/pkg/errors"
+ "github.com/apache/servicecomb-service-center/pkg/gopool"
"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"
@@ -37,6 +38,18 @@ import (
"time"
)
+type ServiceDetailOpt struct {
+ domainProject string
+ service *pb.MicroService
+ countOnly bool
+ options []string
+}
+
+type GetInstanceCountByDomainResponse struct {
+ err error
+ countByDomain int64
+}
+
// service
func capRegisterData(ctx context.Context, request *pb.CreateServiceRequest) (
[]registry.PluginOp, []registry.CompareOp, []registry.PluginOp, error) {
@@ -392,3 +405,253 @@ func revokeInstance(ctx context.Context, domainProject
string, serviceID string,
}
return nil
}
+
+// governServiceCtrl util
+func getServiceAllVersions(ctx context.Context, serviceKey
*pb.MicroServiceKey) ([]string, error) {
+ var versions []string
+
+ copyKey := *serviceKey
+ copyKey.Version = ""
+ key := GenerateServiceIndexKey(©Key)
+
+ opts := append(serviceUtil.FromContext(ctx),
+ registry.WithStrKey(key),
+ registry.WithPrefix())
+
+ resp, err := backend.Store().ServiceIndex().Search(ctx, opts...)
+ if err != nil {
+ return nil, err
+ }
+ if resp == nil || len(resp.Kvs) == 0 {
+ return versions, nil
+ }
+ for _, kv := range resp.Kvs {
+ key := GetInfoFromSvcIndexKV(kv.Key)
+ versions = append(versions, key.Version)
+ }
+ return versions, nil
+}
+
+func getServiceDetailUtil(ctx context.Context, serviceDetailOpt
ServiceDetailOpt) (*pb.ServiceDetail, error) {
+ serviceID := serviceDetailOpt.service.ServiceId
+ options := serviceDetailOpt.options
+ domainProject := serviceDetailOpt.domainProject
+ serviceDetail := new(pb.ServiceDetail)
+ if serviceDetailOpt.countOnly {
+ serviceDetail.Statics = new(pb.Statistics)
+ }
+
+ for _, opt := range options {
+ expr := opt
+ switch expr {
+ case "tags":
+ tags, err := serviceUtil.GetTagsUtils(ctx,
domainProject, serviceID)
+ if err != nil {
+ log.Errorf(err, "get service[%s]'s all tags
failed", serviceID)
+ return nil, err
+ }
+ serviceDetail.Tags = tags
+ case "rules":
+ rules, err := serviceUtil.GetRulesUtil(ctx,
domainProject, serviceID)
+ if err != nil {
+ log.Errorf(err, "get service[%s]'s all rules
failed", serviceID)
+ return nil, err
+ }
+ for _, rule := range rules {
+ rule.Timestamp = rule.ModTimestamp
+ }
+ serviceDetail.Rules = rules
+ case "instances":
+ if serviceDetailOpt.countOnly {
+ instanceCount, err :=
serviceUtil.GetInstanceCountOfOneService(ctx, domainProject, serviceID)
+ if err != nil {
+ log.Errorf(err, "get number of
service[%s]'s instances failed", serviceID)
+ return nil, err
+ }
+ serviceDetail.Statics.Instances =
&pb.StInstance{
+ Count: instanceCount}
+ continue
+ }
+ instances, err :=
serviceUtil.GetAllInstancesOfOneService(ctx, domainProject, serviceID)
+ if err != nil {
+ log.Errorf(err, "get service[%s]'s all
instances failed", serviceID)
+ return nil, err
+ }
+ serviceDetail.Instances = instances
+ case "schemas":
+ schemas, err := getSchemaInfoUtil(ctx, domainProject,
serviceID)
+ if err != nil {
+ log.Errorf(err, "get service[%s]'s all schemas
failed", serviceID)
+ return nil, err
+ }
+ serviceDetail.SchemaInfos = schemas
+ case "dependencies":
+ service := serviceDetailOpt.service
+ dr := serviceUtil.NewDependencyRelation(ctx,
domainProject, service, service)
+ consumers, err := dr.GetDependencyConsumers(
+ serviceUtil.WithoutSelfDependency(),
+ serviceUtil.WithSameDomainProject())
+ if err != nil {
+ log.Errorf(err, "get service[%s][%s/%s/%s/%s]'s
all consumers failed",
+ service.ServiceId, service.Environment,
service.AppId, service.ServiceName, service.Version)
+ return nil, err
+ }
+ providers, err := dr.GetDependencyProviders(
+ serviceUtil.WithoutSelfDependency(),
+ serviceUtil.WithSameDomainProject())
+ if err != nil {
+ log.Errorf(err, "get service[%s][%s/%s/%s/%s]'s
all providers failed",
+ service.ServiceId, service.Environment,
service.AppId, service.ServiceName, service.Version)
+ return nil, err
+ }
+
+ serviceDetail.Consumers = consumers
+ serviceDetail.Providers = providers
+ case "":
+ continue
+ default:
+ log.Errorf(nil, "request option[%s] is invalid", opt)
+ }
+ }
+ return serviceDetail, nil
+}
+
+func getSchemaInfoUtil(ctx context.Context, domainProject string, serviceID
string) ([]*pb.Schema, error) {
+ key := apt.GenerateServiceSchemaKey(domainProject, serviceID, "")
+
+ resp, err := backend.Store().Schema().Search(ctx,
+ registry.WithStrKey(key),
+ registry.WithPrefix())
+ if err != nil {
+ log.Errorf(err, "get service[%s]'s schemas failed", serviceID)
+ return make([]*pb.Schema, 0), err
+ }
+ schemas := make([]*pb.Schema, 0, len(resp.Kvs))
+ for _, kv := range resp.Kvs {
+ schemaInfo := &pb.Schema{}
+ schemaInfo.Schema =
util.BytesToStringWithNoCopy(kv.Value.([]byte))
+ schemaInfo.SchemaId =
util.BytesToStringWithNoCopy(kv.Key[len(key):])
+ schemas = append(schemas, schemaInfo)
+ }
+ return schemas, nil
+}
+
+func statistics(ctx context.Context, withShared bool) (*pb.Statistics, error) {
+ result := &pb.Statistics{
+ Services: &pb.StService{},
+ Instances: &pb.StInstance{},
+ Apps: &pb.StApp{},
+ }
+ domainProject := util.ParseDomainProject(ctx)
+ opts := serviceUtil.FromContext(ctx)
+
+ // services
+ key := apt.GetServiceIndexRootKey(domainProject) + "/"
+ svcOpts := append(opts,
+ registry.WithStrKey(key),
+ registry.WithPrefix())
+ respSvc, err := backend.Store().ServiceIndex().Search(ctx, svcOpts...)
+ if err != nil {
+ return nil, err
+ }
+
+ app := make(map[string]struct{}, respSvc.Count)
+ svcWithNonVersion := make(map[string]struct{}, respSvc.Count)
+ svcIDToNonVerKey := make(map[string]string, respSvc.Count)
+ for _, kv := range respSvc.Kvs {
+ key := apt.GetInfoFromSvcIndexKV(kv.Key)
+ if !withShared && apt.IsShared(key) {
+ continue
+ }
+ if _, ok := app[key.AppId]; !ok {
+ app[key.AppId] = struct{}{}
+ }
+
+ key.Version = ""
+ svcWithNonVersionKey := apt.GenerateServiceIndexKey(key)
+ if _, ok := svcWithNonVersion[svcWithNonVersionKey]; !ok {
+ svcWithNonVersion[svcWithNonVersionKey] = struct{}{}
+ }
+ svcIDToNonVerKey[kv.Value.(string)] = svcWithNonVersionKey
+ }
+
+ result.Services.Count = int64(len(svcWithNonVersion))
+ result.Apps.Count = int64(len(app))
+
+ respGetInstanceCountByDomain := make(chan
GetInstanceCountByDomainResponse, 1)
+ gopool.Go(func(_ context.Context) {
+ getInstanceCountByDomain(ctx, svcIDToNonVerKey,
respGetInstanceCountByDomain)
+ })
+
+ // instance
+ key = apt.GetInstanceRootKey(domainProject) + "/"
+ instOpts := append(opts,
+ registry.WithStrKey(key),
+ registry.WithPrefix(),
+ registry.WithKeyOnly())
+ respIns, err := backend.Store().Instance().Search(ctx, instOpts...)
+ if err != nil {
+ return nil, err
+ }
+
+ onlineServices := make(map[string]struct{}, respSvc.Count)
+ for _, kv := range respIns.Kvs {
+ serviceID, _, _ := apt.GetInfoFromInstKV(kv.Key)
+ key, ok := svcIDToNonVerKey[serviceID]
+ if !ok {
+ continue
+ }
+ result.Instances.Count++
+ if _, ok := onlineServices[key]; !ok {
+ onlineServices[key] = struct{}{}
+ }
+ }
+ result.Services.OnlineCount = int64(len(onlineServices))
+
+ data := <-respGetInstanceCountByDomain
+ close(respGetInstanceCountByDomain)
+ if data.err != nil {
+ return nil, data.err
+ }
+ result.Instances.CountByDomain = data.countByDomain
+ return result, nil
+}
+
+func getInstanceCountByDomain(ctx context.Context, svcIDToNonVerKey
map[string]string, resp chan GetInstanceCountByDomainResponse) {
+ domainID := util.ParseDomain(ctx)
+ key := apt.GetInstanceRootKey(domainID) + "/"
+ instOpts := append([]registry.PluginOpOption{},
+ registry.WithStrKey(key),
+ registry.WithPrefix(),
+ registry.WithKeyOnly())
+ respIns, err := backend.Store().Instance().Search(ctx, instOpts...)
+ ret := GetInstanceCountByDomainResponse{
+ err: err,
+ }
+
+ if err != nil {
+ log.Errorf(err, "get number of instances by domain[%s]",
domainID)
+ } else {
+ for _, kv := range respIns.Kvs {
+ serviceID, _, _ := apt.GetInfoFromInstKV(kv.Key)
+ _, ok := svcIDToNonVerKey[serviceID]
+ if !ok {
+ continue
+ }
+ ret.countByDomain++
+ }
+ }
+
+ resp <- ret
+}
+
+// dep util
+func toDependencyFilterOptions(in *pb.GetDependenciesRequest) (opts
[]serviceUtil.DependencyRelationFilterOption) {
+ if in.SameDomain {
+ opts = append(opts, serviceUtil.WithSameDomainProject())
+ }
+ if in.NoSelf {
+ opts = append(opts, serviceUtil.WithoutSelfDependency())
+ }
+ return opts
+}
diff --git a/datasource/ms.go b/datasource/ms.go
index 52cab0b..424530b 100644
--- a/datasource/ms.go
+++ b/datasource/ms.go
@@ -28,6 +28,11 @@ type MetadataManager interface {
RegisterService(ctx context.Context, request *pb.CreateServiceRequest)
(*pb.CreateServiceResponse, error)
GetServices(ctx context.Context, request *pb.GetServicesRequest)
(*pb.GetServicesResponse, error)
GetService(ctx context.Context, request *pb.GetServiceRequest)
(*pb.GetServiceResponse, error)
+
+ GetServiceDetail(ctx context.Context, request *pb.GetServiceRequest)
(*pb.GetServiceDetailResponse, error)
+ GetServicesInfo(ctx context.Context, request
*pb.GetServicesInfoRequest) (*pb.GetServicesInfoResponse, error)
+ GetApplications(ctx context.Context, request *pb.GetAppsRequest)
(*pb.GetAppsResponse, error)
+
ExistService(ctx context.Context, request *pb.GetExistenceRequest)
(*pb.GetExistenceResponse, error)
UpdateService(ctx context.Context, request
*pb.UpdateServicePropsRequest) (*pb.UpdateServicePropsResponse, error)
UnregisterService(ctx context.Context, request
*pb.DeleteServiceRequest) (*pb.DeleteServiceResponse, error)