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 aa698f7 [SCB 2094] implement datasource/ms tag interface (#719)
aa698f7 is described below
commit aa698f76210866b7582cbee84e0c564ddee1f6d4
Author: robotLJW <[email protected]>
AuthorDate: Mon Oct 19 14:09:03 2020 +0800
[SCB 2094] implement datasource/ms tag interface (#719)
1. implement interface
+ GetTags
+ UpdateTag
+ DeleteTags
2. finish unit test relative to the instance above
---
datasource/etcd/ms.go | 154 +++++++++++++++++++++-
datasource/etcd/ms_test.go | 315 +++++++++++++++++++++++++++++++++++++++++++++
datasource/ms.go | 8 +-
3 files changed, 467 insertions(+), 10 deletions(-)
diff --git a/datasource/etcd/ms.go b/datasource/etcd/ms.go
index c84ef79..cabf9f1 100644
--- a/datasource/etcd/ms.go
+++ b/datasource/etcd/ms.go
@@ -1304,16 +1304,158 @@ func (ds *DataSource) AddTags(ctx context.Context,
request *pb.AddServiceTagsReq
}, nil
}
-func (ds *DataSource) GetTag() {
- panic("implement me")
+func (ds *DataSource) GetTags(ctx context.Context, request
*pb.GetServiceTagsRequest) (*pb.GetServiceTagsResponse, error) {
+ var err error
+ domainProject := util.ParseDomainProject(ctx)
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(err, "get service[%s]'s tags failed, service does
not exist", request.ServiceId)
+ return &pb.GetServiceTagsResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+ tags, err := serviceUtil.GetTagsUtils(ctx, domainProject,
request.ServiceId)
+ if err != nil {
+ log.Errorf(err, "get service[%s]'s tags failed, get tags
failed", request.ServiceId)
+ return &pb.GetServiceTagsResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ return &pb.GetServiceTagsResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Get
service tags successfully."),
+ Tags: tags,
+ }, nil
}
-func (ds *DataSource) UpdateTag() {
- panic("implement me")
+func (ds *DataSource) UpdateTag(ctx context.Context, request
*pb.UpdateServiceTagRequest) (*pb.UpdateServiceTagResponse, error) {
+ var err error
+ remoteIP := util.GetIPFromContext(ctx)
+ tagFlag := util.StringJoin([]string{request.Key, request.Value}, "/")
+ domainProject := util.ParseDomainProject(ctx)
+
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(err, "update service[%s]'s tag[%s] failed, service
does not exist, operator: %s",
+ request.ServiceId, tagFlag, remoteIP)
+ return &pb.UpdateServiceTagResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ tags, err := serviceUtil.GetTagsUtils(ctx, domainProject,
request.ServiceId)
+ if err != nil {
+ log.Errorf(err, "update service[%s]'s tag[%s] failed, get tag
failed, operator: %s",
+ request.ServiceId, tagFlag, remoteIP)
+ return &pb.UpdateServiceTagResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ //check if the tag exists
+ if _, ok := tags[request.Key]; !ok {
+ log.Errorf(nil, "update service[%s]'s tag[%s] failed, tag does
not exist, operator: %s",
+ request.ServiceId, tagFlag, remoteIP)
+ return &pb.UpdateServiceTagResponse{
+ Response: proto.CreateResponse(scerr.ErrTagNotExists,
"Tag does not exist, please add one first."),
+ }, nil
+ }
+
+ copyTags := make(map[string]string, len(tags))
+ for k, v := range tags {
+ copyTags[k] = v
+ }
+ copyTags[request.Key] = request.Value
+
+ checkErr := serviceUtil.AddTagIntoETCD(ctx, domainProject,
request.ServiceId, copyTags)
+ if checkErr != nil {
+ log.Errorf(checkErr, "update service[%s]'s tag[%s] failed,
operator: %s", request.ServiceId, tagFlag, remoteIP)
+ resp := &pb.UpdateServiceTagResponse{
+ Response: proto.CreateResponseWithSCErr(checkErr),
+ }
+ if checkErr.InternalError() {
+ return resp, checkErr
+ }
+ return resp, nil
+ }
+
+ log.Infof("update service[%s]'s tag[%s] successfully, operator: %s",
request.ServiceId, tagFlag, remoteIP)
+ return &pb.UpdateServiceTagResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Update
service tag success."),
+ }, nil
}
-func (ds *DataSource) DeleteTag() {
- panic("implement me")
+func (ds *DataSource) DeleteTags(ctx context.Context, request
*pb.DeleteServiceTagsRequest) (*pb.DeleteServiceTagsResponse, error) {
+ remoteIP := util.GetIPFromContext(ctx)
+ domainProject := util.ParseDomainProject(ctx)
+
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(nil, "delete service[%s]'s tags %v failed, service
does not exist, operator: %s",
+ request.ServiceId, request.Keys, remoteIP)
+ return &pb.DeleteServiceTagsResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ tags, err := serviceUtil.GetTagsUtils(ctx, domainProject,
request.ServiceId)
+ if err != nil {
+ log.Errorf(err, "delete service[%s]'s tags %v failed, get
service tags failed, operator: %s",
+ request.ServiceId, request.Keys, remoteIP)
+ return &pb.DeleteServiceTagsResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ copyTags := make(map[string]string, len(tags))
+ for k, v := range tags {
+ copyTags[k] = v
+ }
+ for _, key := range request.Keys {
+ if _, ok := copyTags[key]; !ok {
+ log.Errorf(nil, "delete service[%s]'s tags %v failed,
tag[%s] does not exist, operator: %s",
+ request.ServiceId, request.Keys, key, remoteIP)
+ return &pb.DeleteServiceTagsResponse{
+ Response:
proto.CreateResponse(scerr.ErrTagNotExists, "Delete tags failed for this key
"+key+" does not exist."),
+ }, nil
+ }
+ delete(copyTags, key)
+ }
+
+ // the capacity of tags may be 0
+ data, err := json.Marshal(copyTags)
+ if err != nil {
+ log.Errorf(err, "delete service[%s]'s tags %v failed, marshall
service tags failed, operator: %s",
+ request.ServiceId, request.Keys, remoteIP)
+ return &pb.DeleteServiceTagsResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ key := apt.GenerateServiceTagKey(domainProject, request.ServiceId)
+
+ resp, err := backend.Registry().TxnWithCmp(ctx,
+ []registry.PluginOp{registry.OpPut(registry.WithStrKey(key),
registry.WithValue(data))},
+ []registry.CompareOp{registry.OpCmp(
+
registry.CmpVer(util.StringToBytesWithNoCopy(apt.GenerateServiceKey(domainProject,
request.ServiceId))),
+ registry.CmpNotEqual, 0)},
+ nil)
+ if err != nil {
+ log.Errorf(err, "delete service[%s]'s tags %v failed, operator:
%s",
+ request.ServiceId, request.Keys, remoteIP)
+ return &pb.DeleteServiceTagsResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, err.Error()),
+ }, err
+ }
+ if !resp.Succeeded {
+ log.Errorf(err, "delete service[%s]'s tags %v failed, service
does not exist, operator: %s",
+ request.ServiceId, request.Keys, remoteIP)
+ return &pb.DeleteServiceTagsResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ log.Infof("delete service[%s]'s tags %v successfully, operator: %s",
request.ServiceId, request.Keys, remoteIP)
+ return &pb.DeleteServiceTagsResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Delete
service tags successfully."),
+ }, nil
}
func (ds *DataSource) AddRule(ctx context.Context, request
*pb.AddServiceRulesRequest) (
diff --git a/datasource/etcd/ms_test.go b/datasource/etcd/ms_test.go
index 948be0a..49e5d83 100644
--- a/datasource/etcd/ms_test.go
+++ b/datasource/etcd/ms_test.go
@@ -3688,3 +3688,318 @@ func TestTags_Add(t *testing.T) {
assert.Equal(t, scerr.ErrNotEnoughQuota,
resp.Response.GetCode())
})
}
+
+func TestTags_Get(t *testing.T) {
+ var serviceId string
+ 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("create service and add tags", func(t *testing.T) {
+ svc := &pb.MicroService{
+ AppId: "get_tag_group_ms",
+ ServiceName: "get_tag_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ }
+ resp, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: svc,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ serviceId = resp.ServiceId
+
+ log.Info("add tags should be passed")
+ respAddTags, err := datasource.Instance().AddTags(getContext(),
&pb.AddServiceTagsRequest{
+ ServiceId: serviceId,
+ Tags: map[string]string{
+ "a": "test",
+ "b": "b",
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddTags.Response.GetCode())
+ })
+
+ t.Run("the request is invalid", func(t *testing.T) {
+ log.Info("service does not exists")
+ resp, err := datasource.Instance().GetTags(getContext(),
&pb.GetServiceTagsRequest{
+ ServiceId: "noThisService",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("service's id is empty")
+ resp, err = datasource.Instance().GetTags(getContext(),
&pb.GetServiceTagsRequest{
+ ServiceId: "",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("service's id is invalid")
+ resp, err = datasource.Instance().GetTags(getContext(),
&pb.GetServiceTagsRequest{
+ ServiceId: strings.Repeat("x", 65),
+ })
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+ })
+
+ t.Run("the request is valid", func(t *testing.T) {
+ resp, err := datasource.Instance().GetTags(getContext(),
&pb.GetServiceTagsRequest{
+ ServiceId: serviceId,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ assert.Equal(t, "test", resp.Tags["a"])
+ })
+}
+
+func TestTag_Update(t *testing.T) {
+ var serviceId string
+ 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("add service and add tags", func(t *testing.T) {
+ svc := &pb.MicroService{
+ AppId: "update_tag_group_ms",
+ ServiceName: "update_tag_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ }
+ resp, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: svc,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ serviceId = resp.ServiceId
+
+ log.Info("add tags")
+ respAddTags, err := datasource.Instance().AddTags(getContext(),
&pb.AddServiceTagsRequest{
+ ServiceId: serviceId,
+ Tags: map[string]string{
+ "a": "test",
+ "b": "b",
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddTags.Response.GetCode())
+ })
+
+ t.Run("the request is invalid", func(t *testing.T) {
+
+ log.Info("service does not exists")
+ resp, err := datasource.Instance().UpdateTag(getContext(),
&pb.UpdateServiceTagRequest{
+ ServiceId: "noneservice",
+ Key: "a",
+ Value: "update",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("tag key does not exist")
+ resp, err = datasource.Instance().UpdateTag(getContext(),
&pb.UpdateServiceTagRequest{
+ ServiceId: serviceId,
+ Key: "notexisttag",
+ Value: "update",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrTagNotExists, resp.Response.GetCode())
+
+ log.Info("tag key is invalid")
+ resp, err = datasource.Instance().UpdateTag(getContext(),
&pb.UpdateServiceTagRequest{
+ ServiceId: serviceId,
+ Key: strings.Repeat("x", 65),
+ Value: "v",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrTagNotExists, resp.Response.GetCode())
+ })
+
+ t.Run("the request is valid", func(t *testing.T) {
+ resp, err := datasource.Instance().UpdateTag(getContext(),
&pb.UpdateServiceTagRequest{
+ ServiceId: serviceId,
+ Key: "a",
+ Value: "update",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+
+ t.Run("find instance, contain tag", func(t *testing.T) {
+ log.Info("create consumer")
+ resp, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "find_inst_tag_group_ms",
+ ServiceName: "find_inst_tag_consumer_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ consumerId := resp.ServiceId
+
+ log.Info("create provider")
+ resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "find_inst_tag_group_ms",
+ ServiceName: "find_inst_tag_provider_ms",
+ Version: "1.0.1",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ providerId := resp.ServiceId
+
+ log.Info("tag the provider")
+ addTagsResp, err := datasource.Instance().AddTags(getContext(),
&pb.AddServiceTagsRequest{
+ ServiceId: providerId,
+ Tags: map[string]string{"filter_tag": "filter"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
addTagsResp.Response.GetCode())
+
+ log.Info("add instance to provider")
+ instanceResp, err :=
datasource.Instance().RegisterInstance(getContext(),
&pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: providerId,
+ Endpoints: []string{
+
"findInstanceForTagFilter:127.0.0.1:8080",
+ },
+ HostName: "UT-HOST",
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
instanceResp.Response.GetCode())
+
+ log.Info("find instance")
+ findResp, err :=
datasource.Instance().FindInstances(getContext(), &pb.FindInstancesRequest{
+ ConsumerServiceId: consumerId,
+ AppId: "find_inst_tag_group_ms",
+ ServiceName: "find_inst_tag_provider_ms",
+ VersionRule: "1.0.0+",
+ Tags: []string{"not-exist-tag"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
findResp.Response.GetCode())
+ assert.Equal(t, 0, len(findResp.Instances))
+
+ findResp, err =
datasource.Instance().FindInstances(getContext(), &pb.FindInstancesRequest{
+ ConsumerServiceId: consumerId,
+ AppId: "find_inst_tag_group_ms",
+ ServiceName: "find_inst_tag_provider_ms",
+ VersionRule: "1.0.0+",
+ Tags: []string{"filter_tag"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
findResp.Response.GetCode())
+ assert.Equal(t, instanceResp.InstanceId,
findResp.Instances[0].InstanceId)
+
+ // no add rules
+
+ log.Info("add tags")
+ addTagsResp, err = datasource.Instance().AddTags(getContext(),
&pb.AddServiceTagsRequest{
+ ServiceId: consumerId,
+ Tags: map[string]string{"consumer_tag": "filter"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
addTagsResp.Response.GetCode())
+
+ findResp, err =
datasource.Instance().FindInstances(getContext(), &pb.FindInstancesRequest{
+ ConsumerServiceId: consumerId,
+ AppId: "find_inst_tag_group_ms",
+ ServiceName: "find_inst_tag_provider_ms",
+ VersionRule: "1.0.0+",
+ Tags: []string{"filter_tag"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
findResp.Response.GetCode())
+ })
+}
+
+func TestTags_Delete(t *testing.T) {
+ var serviceId string
+ 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("create service and add tags", func(t *testing.T) {
+ resp, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "delete_tag_group_ms",
+ ServiceName: "delete_tag_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ serviceId = resp.ServiceId
+
+ respAddTages, err :=
datasource.Instance().AddTags(getContext(), &pb.AddServiceTagsRequest{
+ ServiceId: serviceId,
+ Tags: map[string]string{
+ "a": "test",
+ "b": "b",
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddTages.Response.GetCode())
+ })
+
+ t.Run("the request is invalid", func(t *testing.T) {
+ log.Info("service does not exits")
+ resp, err := datasource.Instance().DeleteTags(getContext(),
&pb.DeleteServiceTagsRequest{
+ ServiceId: "noneservice",
+ Keys: []string{"a", "b"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("tag key does not exits")
+ resp, err = datasource.Instance().DeleteTags(getContext(),
&pb.DeleteServiceTagsRequest{
+ ServiceId: serviceId,
+ Keys: []string{"c"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrTagNotExists, resp.Response.GetCode())
+ })
+
+ t.Run("the request is valid", func(t *testing.T) {
+ resp, err := datasource.Instance().DeleteTags(getContext(),
&pb.DeleteServiceTagsRequest{
+ ServiceId: serviceId,
+ Keys: []string{"a", "b"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ respGetTags, err := datasource.Instance().GetTags(getContext(),
&pb.GetServiceTagsRequest{
+ ServiceId: serviceId,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ assert.Equal(t, "", respGetTags.Tags["a"])
+ })
+}
diff --git a/datasource/ms.go b/datasource/ms.go
index c7bd219..52cab0b 100644
--- a/datasource/ms.go
+++ b/datasource/ms.go
@@ -55,10 +55,10 @@ type MetadataManager interface {
GetAllSchemas(ctx context.Context, request *pb.GetAllSchemaRequest)
(*pb.GetAllSchemaResponse, error)
DeleteSchema(ctx context.Context, request *pb.DeleteSchemaRequest)
(*pb.DeleteSchemaResponse, error)
- AddTags(ctx context.Context, in *pb.AddServiceTagsRequest)
(*pb.AddServiceTagsResponse, error)
- GetTag()
- UpdateTag()
- DeleteTag()
+ AddTags(ctx context.Context, request *pb.AddServiceTagsRequest)
(*pb.AddServiceTagsResponse, error)
+ GetTags(ctx context.Context, request *pb.GetServiceTagsRequest)
(*pb.GetServiceTagsResponse, error)
+ UpdateTag(ctx context.Context, request *pb.UpdateServiceTagRequest)
(*pb.UpdateServiceTagResponse, error)
+ DeleteTags(ctx context.Context, request *pb.DeleteServiceTagsRequest)
(*pb.DeleteServiceTagsResponse, error)
AddRule(ctx context.Context, request *pb.AddServiceRulesRequest)
(*pb.AddServiceRulesResponse, error)
GetRule(ctx context.Context, request *pb.GetServiceRulesRequest)
(*pb.GetServiceRulesResponse, error)