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 4251967 [SCB 2094] implement ms module interface (#710) (#718)
4251967 is described below
commit 425196780132d0505da7f343b3b969a856fc8443
Author: popozy <[email protected]>
AuthorDate: Sun Oct 18 15:31:24 2020 +0800
[SCB 2094] implement ms module interface (#710) (#718)
1. implement service/ms/etcd interface
+ GetSchema
+ GetAllSchema
+ DeleteSchema
+ AddRule
+ GetRule
+ DeleteRule
+ UpdateRule
2. finish unit test relative to the instance above
3. add datasource to unit shell script
---
datasource/etcd/dep_test.go | 268 ++++++------
datasource/etcd/etcd_suite_test.go | 6 +
datasource/etcd/ms.go | 697 ++++++++++++++++++++++++++-----
datasource/etcd/ms_test.go | 825 ++++++++++++++++++++++++++++++++++++-
datasource/ms.go | 13 +-
scripts/ut_test_in_docker.sh | 1 +
6 files changed, 1559 insertions(+), 251 deletions(-)
diff --git a/datasource/etcd/dep_test.go b/datasource/etcd/dep_test.go
index c7659ee..f2d07e6 100644
--- a/datasource/etcd/dep_test.go
+++ b/datasource/etcd/dep_test.go
@@ -49,10 +49,10 @@ func Test_Creat(t *testing.T) {
consumerId3 string
)
t.Run("should be passed", func(t *testing.T) {
- resp, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ resp, err :=
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "create_dep_group",
- ServiceName: "create_dep_consumer",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_consumer",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -63,10 +63,10 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
consumerId1 = resp.ServiceId
- resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ resp, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "create_dep_group",
- ServiceName: "create_dep_consumer_all",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_consumer_all",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -77,11 +77,11 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
consumerId3 = resp.ServiceId
- resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ resp, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
Environment: pb.ENV_PROD,
- AppId: "create_dep_group",
- ServiceName: "create_dep_consumer",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_consumer",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -92,10 +92,10 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
consumerId2 = resp.ServiceId
- resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ resp, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "create_dep_group",
- ServiceName: "create_dep_provider",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_provider",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -105,10 +105,10 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
- resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ resp, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "create_dep_group",
- ServiceName: "create_dep_provider",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_provider",
Version: "1.0.1",
Level: "FRONT",
Status: pb.MS_UP,
@@ -118,11 +118,11 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
- resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ resp, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
Environment: pb.ENV_PROD,
- AppId: "create_dep_group",
- ServiceName: "create_dep_provider",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_provider",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -135,20 +135,20 @@ func Test_Creat(t *testing.T) {
t.Run("when request is invalid, should be failed", func(t *testing.T) {
consumer := &pb.MicroServiceKey{
- AppId: "create_dep_group",
- ServiceName: "create_dep_consumer",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_consumer",
Version: "1.0.0",
}
providers := []*pb.MicroServiceKey{
{
- AppId: "create_dep_group",
- ServiceName: "create_dep_provider",
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_provider",
Version: "1.0.0",
},
}
// consumer does not exist
- resp, err :=
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err :=
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
@@ -165,14 +165,14 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
// provider version is invalid
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"create_dep_group",
- ServiceName:
"create_dep_provider",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_provider",
Version:
"1.0.32768",
},
},
@@ -184,12 +184,12 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
// consumer version is invalid
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
- AppId: "create_dep_group",
- ServiceName:
"create_dep_consumer",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_consumer",
Version: "1.0.0+",
},
Providers: providers,
@@ -200,12 +200,12 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
- AppId: "create_dep_group",
- ServiceName:
"create_dep_consumer",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_consumer",
Version: "1.0.0-1.0.1",
},
Providers: providers,
@@ -216,12 +216,12 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
- AppId: "create_dep_group",
- ServiceName:
"create_dep_consumer",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_consumer",
Version: "latest",
},
Providers: providers,
@@ -232,12 +232,12 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
- AppId: "create_dep_group",
- ServiceName:
"create_dep_consumer",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_consumer",
Version: "",
},
Providers: providers,
@@ -248,11 +248,11 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
- AppId: "create_dep_group",
+ AppId:
"dep_create_dep_group",
ServiceName: "*",
Version: "1.0.0",
},
@@ -265,14 +265,14 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
// provider app is invalid
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
AppId: "*",
- ServiceName:
"service_name_provider",
+ ServiceName:
"dep_service_name_provider",
Version: "2.0.0",
},
},
@@ -284,13 +284,13 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
// provider serviceName is invalid
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"service_group_provider",
+ AppId:
"dep_service_group_provider",
ServiceName: "-",
Version: "2.0.0",
},
@@ -303,14 +303,14 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
// provider version is invalid
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"service_group_provider",
- ServiceName:
"service_name_provider",
+ AppId:
"dep_service_group_provider",
+ ServiceName:
"dep_service_name_provider",
Version: "",
},
},
@@ -322,15 +322,15 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
// provider in diff env
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
Environment:
pb.ENV_PROD,
- AppId:
"service_group_provider",
- ServiceName:
"service_name_provider",
+ AppId:
"dep_service_group_provider",
+ ServiceName:
"dep_service_name_provider",
Version: "latest",
},
},
@@ -343,14 +343,14 @@ func Test_Creat(t *testing.T) {
// consumer in diff env
consumer.Environment = pb.ENV_PROD
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"service_group_provider",
- ServiceName:
"service_name_provider",
+ AppId:
"dep_service_group_provider",
+ ServiceName:
"dep_service_name_provider",
Version: "latest",
},
},
@@ -363,7 +363,7 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respCon, err :=
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respCon, err :=
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, respCon)
@@ -371,7 +371,7 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respCon.Response.GetCode())
assert.Equal(t, 0, len(respCon.Providers))
- respCon, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respCon, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId2,
})
assert.NotNil(t, respCon)
@@ -384,20 +384,20 @@ func Test_Creat(t *testing.T) {
for i := 0; i < 101; i++ {
deps = append(deps, &pb.ConsumerDependency{
Consumer: &pb.MicroServiceKey{
- AppId: "create_dep_group",
- ServiceName: "create_dep_consumer" +
strconv.Itoa(i),
+ AppId: "dep_create_dep_group",
+ ServiceName: "dep_create_dep_consumer"
+ strconv.Itoa(i),
Version: "1.0.0",
},
Providers: []*pb.MicroServiceKey{
{
- AppId:
"service_group_provider",
- ServiceName:
"service_name_provider",
+ AppId:
"dep_service_group_provider",
+ ServiceName:
"dep_service_name_provider",
Version: "latest",
},
},
})
}
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: deps,
})
assert.NotNil(t, resp)
@@ -407,20 +407,20 @@ func Test_Creat(t *testing.T) {
t.Run("when request is valid, should be passed", func(t *testing.T) {
consumer := &pb.MicroServiceKey{
- ServiceName: "create_dep_consumer",
- AppId: "create_dep_group",
+ ServiceName: "dep_create_dep_consumer",
+ AppId: "dep_create_dep_group",
Version: "1.0.0",
}
// add latest
- resp, err :=
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err :=
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"create_dep_group",
- ServiceName:
"create_dep_provider",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_provider",
Version: "latest",
},
},
@@ -433,7 +433,7 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respPro, err :=
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respPro, err :=
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, respPro)
@@ -441,14 +441,14 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respPro.Response.GetCode())
// add 1.0.0+
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"create_dep_group",
- ServiceName:
"create_dep_provider",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_provider",
Version: "1.0.0+",
},
},
@@ -459,7 +459,7 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
- respPro, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respPro, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, respPro)
@@ -467,12 +467,12 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respPro.Response.GetCode())
// add *
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
- ServiceName:
"create_dep_consumer_all",
- AppId: "create_dep_group",
+ ServiceName:
"dep_create_dep_consumer_all",
+ AppId:
"dep_create_dep_group",
Version: "1.0.0",
},
Providers: []*pb.MicroServiceKey{
@@ -487,7 +487,7 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
- respPro, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respPro, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId3,
})
assert.NotNil(t, respPro)
@@ -496,12 +496,12 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, 0, len(respPro.Providers))
// clean all
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: &pb.MicroServiceKey{
- ServiceName:
"create_dep_consumer_all",
- AppId: "create_dep_group",
+ ServiceName:
"dep_create_dep_consumer_all",
+ AppId:
"dep_create_dep_group",
Version: "1.0.0",
},
Providers: nil,
@@ -513,14 +513,14 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
// add multiple providers
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"create_dep_group",
- ServiceName:
"create_dep_provider",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_provider",
Version: "1.0.0",
},
{
@@ -535,14 +535,14 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
// add 1.0.0-2.0.0 to override *
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"create_dep_group",
- ServiceName:
"create_dep_provider",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_provider",
Version:
"1.0.0-1.0.1",
},
},
@@ -555,7 +555,7 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respPro, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respPro, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, respPro)
@@ -564,14 +564,14 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, "1.0.0", respPro.Providers[0].Version)
// add not override
- respAdd, err :=
datasource.Instance().AddDependency(getContext(), &pb.AddDependenciesRequest{
+ respAdd, err :=
datasource.Instance().AddDependency(depGetContext(), &pb.AddDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
Providers: []*pb.MicroServiceKey{
{
- AppId:
"create_dep_group",
- ServiceName:
"create_dep_provider",
+ AppId:
"dep_create_dep_group",
+ ServiceName:
"dep_create_dep_provider",
Version:
"1.0.0-3.0.0",
},
},
@@ -584,7 +584,7 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respPro, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respPro, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, respPro)
@@ -592,7 +592,7 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respPro.Response.GetCode())
// add provider is empty
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
@@ -604,7 +604,7 @@ func Test_Creat(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
- resp, err =
datasource.Instance().CreateDependency(getContext(),
&pb.CreateDependenciesRequest{
+ resp, err =
datasource.Instance().CreateDependency(depGetContext(),
&pb.CreateDependenciesRequest{
Dependencies: []*pb.ConsumerDependency{
{
Consumer: consumer,
@@ -617,7 +617,7 @@ func Test_Creat(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respPro, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respPro, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, respPro)
@@ -646,10 +646,10 @@ func Test_Get(t *testing.T) {
)
t.Run("should be passed", func(t *testing.T) {
- resp, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ resp, err :=
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "get_dep_group",
- ServiceName: "get_dep_consumer",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_consumer",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -660,10 +660,10 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
consumerId1 = resp.ServiceId
- resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ resp, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "get_dep_group",
- ServiceName: "get_dep_provider",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_provider",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -674,10 +674,10 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
providerId1 = resp.ServiceId
- resp, err = datasource.Instance().RegisterService(getContext(),
&pb.CreateServiceRequest{
+ resp, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "get_dep_group",
- ServiceName: "get_dep_provider",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_provider",
Version: "2.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -691,7 +691,7 @@ func Test_Get(t *testing.T) {
t.Run("when request is invalid, should be failed", func(t *testing.T) {
//service id is empty when get provider
- resp, err :=
datasource.Instance().SearchProviderDependency(getContext(),
&pb.GetDependenciesRequest{
+ resp, err :=
datasource.Instance().SearchProviderDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: "",
})
assert.NotNil(t, resp)
@@ -699,7 +699,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
//service does not exist when get provider
- resp, err =
datasource.Instance().SearchProviderDependency(getContext(),
&pb.GetDependenciesRequest{
+ resp, err =
datasource.Instance().SearchProviderDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: "noneservice",
})
assert.NotNil(t, resp)
@@ -707,7 +707,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
//service id is empty when get consumer
- resp, err =
datasource.Instance().SearchProviderDependency(getContext(),
&pb.GetDependenciesRequest{
+ resp, err =
datasource.Instance().SearchProviderDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: "",
})
assert.NotNil(t, resp)
@@ -715,7 +715,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, scerr.ErrInvalidParams, resp.Response.GetCode())
//service does not exist when get consumer
- resp, err =
datasource.Instance().SearchProviderDependency(getContext(),
&pb.GetDependenciesRequest{
+ resp, err =
datasource.Instance().SearchProviderDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: "noneservice",
})
assert.NotNil(t, resp)
@@ -725,7 +725,7 @@ func Test_Get(t *testing.T) {
t.Run("when request is valid, should be passed", func(t *testing.T) {
//get provider
- resp, err :=
datasource.Instance().SearchProviderDependency(getContext(),
&pb.GetDependenciesRequest{
+ resp, err :=
datasource.Instance().SearchProviderDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: providerId1,
})
assert.NotNil(t, resp)
@@ -733,7 +733,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
//get consumer
- resp, err =
datasource.Instance().SearchProviderDependency(getContext(),
&pb.GetDependenciesRequest{
+ resp, err =
datasource.Instance().SearchProviderDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, resp)
@@ -743,10 +743,10 @@ func Test_Get(t *testing.T) {
t.Run("when after finding instance, should created dependencies between
C and P", func(t *testing.T) {
// find provider
- resp, err := datasource.Instance().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ resp, err :=
datasource.Instance().FindInstances(depGetContext(), &pb.FindInstancesRequest{
ConsumerServiceId: consumerId1,
- AppId: "get_dep_group",
- ServiceName: "get_dep_provider",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_provider",
VersionRule: "1.0.0+",
})
assert.NotNil(t, resp)
@@ -756,7 +756,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
// get consumer's deps
- respGetP, err :=
datasource.Instance().SearchProviderDependency(getContext(),
&pb.GetDependenciesRequest{
+ respGetP, err :=
datasource.Instance().SearchProviderDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: providerId1,
})
assert.NotNil(t, respGetP)
@@ -764,7 +764,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respGetP.Response.GetCode())
// get provider's deps
- respGetC, err :=
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respGetC, err :=
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
})
assert.NotNil(t, respGetC)
@@ -772,10 +772,10 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respGetC.Response.GetCode())
// get self deps
- resp, err = datasource.Instance().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ resp, err =
datasource.Instance().FindInstances(depGetContext(), &pb.FindInstancesRequest{
ConsumerServiceId: consumerId1,
- AppId: "get_dep_group",
- ServiceName: "get_dep_consumer",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_consumer",
VersionRule: "1.0.0+",
})
assert.NotNil(t, resp)
@@ -784,7 +784,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respGetC, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respGetC, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: consumerId1,
NoSelf: true,
})
@@ -793,20 +793,20 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respGetC.Response.GetCode())
// find before provider register
- resp, err = datasource.Instance().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ resp, err =
datasource.Instance().FindInstances(depGetContext(), &pb.FindInstancesRequest{
ConsumerServiceId: providerId2,
- AppId: "get_dep_group",
- ServiceName: "get_dep_finder",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_finder",
VersionRule: "1.0.0+",
})
assert.NotNil(t, resp)
assert.NoError(t, err)
assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
- respCreateF, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ respCreateF, err :=
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
- AppId: "get_dep_group",
- ServiceName: "get_dep_finder",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_finder",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -817,10 +817,10 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respCreateF.Response.GetCode())
finder1 := respCreateF.ServiceId
- resp, err = datasource.Instance().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ resp, err =
datasource.Instance().FindInstances(depGetContext(), &pb.FindInstancesRequest{
ConsumerServiceId: providerId2,
- AppId: "get_dep_group",
- ServiceName: "get_dep_finder",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_finder",
VersionRule: "1.0.0+",
})
assert.NotNil(t, resp)
@@ -829,7 +829,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respGetC, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respGetC, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: providerId2,
})
assert.NotNil(t, respGetC)
@@ -839,7 +839,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, finder1, respGetC.Providers[0].ServiceId)
// find after delete micro service
- respDelP, err :=
datasource.Instance().UnregisterService(getContext(), &pb.DeleteServiceRequest{
+ respDelP, err :=
datasource.Instance().UnregisterService(depGetContext(),
&pb.DeleteServiceRequest{
ServiceId: finder1, Force: true,
})
assert.NotNil(t, respDelP)
@@ -848,7 +848,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respGetC, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respGetC, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: providerId2,
})
assert.NotNil(t, respGetC)
@@ -856,11 +856,11 @@ func Test_Get(t *testing.T) {
assert.Equal(t, proto.Response_SUCCESS,
respGetC.Response.GetCode())
assert.Equal(t, 0, len(respGetC.Providers))
- respCreateF, err =
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ respCreateF, err =
datasource.Instance().RegisterService(depGetContext(), &pb.CreateServiceRequest{
Service: &pb.MicroService{
ServiceId: finder1,
- AppId: "get_dep_group",
- ServiceName: "get_dep_finder",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_finder",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -870,10 +870,10 @@ func Test_Get(t *testing.T) {
assert.NoError(t, err)
assert.Equal(t, proto.Response_SUCCESS,
respCreateF.Response.GetCode())
- resp, err = datasource.Instance().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ resp, err =
datasource.Instance().FindInstances(depGetContext(), &pb.FindInstancesRequest{
ConsumerServiceId: providerId2,
- AppId: "get_dep_group",
- ServiceName: "get_dep_finder",
+ AppId: "dep_get_dep_group",
+ ServiceName: "dep_get_dep_finder",
VersionRule: "1.0.0+",
})
assert.NotNil(t, resp)
@@ -882,7 +882,7 @@ func Test_Get(t *testing.T) {
assert.Equal(t, nil, deh.Handle())
- respGetC, err =
datasource.Instance().SearchConsumerDependency(getContext(),
&pb.GetDependenciesRequest{
+ respGetC, err =
datasource.Instance().SearchConsumerDependency(depGetContext(),
&pb.GetDependenciesRequest{
ServiceId: providerId2,
})
assert.NotNil(t, respGetC)
diff --git a/datasource/etcd/etcd_suite_test.go
b/datasource/etcd/etcd_suite_test.go
index 1560595..aa275a0 100644
--- a/datasource/etcd/etcd_suite_test.go
+++ b/datasource/etcd/etcd_suite_test.go
@@ -54,6 +54,12 @@ func getContext() context.Context {
util.CtxNocache, "1")
}
+func depGetContext() context.Context {
+ return util.SetContext(
+ util.SetDomainProject(context.Background(), "new_default",
"new_default"),
+ util.CtxNocache, "1")
+}
+
func TestGrpc(t *testing.T) {
RegisterFailHandler(Fail)
junitReporter := reporters.NewJUnitReporter("model.junit.xml")
diff --git a/datasource/etcd/ms.go b/datasource/etcd/ms.go
index 8b2b41b..c84ef79 100644
--- a/datasource/etcd/ms.go
+++ b/datasource/etcd/ms.go
@@ -3,12 +3,12 @@
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
+ * (the "License"); you may not use this file except request compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
- * Unless required by applicable law or agreed to in writing, software
+ * Unless required by applicable law or agreed to request writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
@@ -30,6 +30,7 @@ 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/quota"
"github.com/apache/servicecomb-service-center/server/plugin/registry"
scerr "github.com/apache/servicecomb-service-center/server/scerror"
@@ -216,9 +217,9 @@ func (ds *DataSource) RegisterInstance(ctx context.Context,
request *pb.Register
if len(instance.InstanceId) > 0 {
// keep alive the lease ttl
// there are two reasons for sending a heartbeat here:
- // 1. in the scenario the instance has been removed,
+ // 1. request the scenario the instance has been removed,
// the cast of registration operation can be reduced.
- // 2. in the self-protection scenario, the instance is unhealthy
+ // 2. request the self-protection scenario, the instance is
unhealthy
// and needs to be re-registered.
resp, err := ds.Heartbeat(ctx, &pb.HeartbeatRequest{ServiceId:
instance.ServiceId,
InstanceId: instance.InstanceId})
@@ -352,61 +353,61 @@ func (ds *DataSource) RegisterInstance(ctx
context.Context, request *pb.Register
}, nil
}
-func (ds *DataSource) GetInstance(ctx context.Context, in
*pb.GetOneInstanceRequest) (
+func (ds *DataSource) GetInstance(ctx context.Context, request
*pb.GetOneInstanceRequest) (
*pb.GetOneInstanceResponse, error) {
domainProject := util.ParseDomainProject(ctx)
service := &pb.MicroService{}
var err error
- if len(in.ConsumerServiceId) > 0 {
- service, err = serviceUtil.GetService(ctx, domainProject,
in.ConsumerServiceId)
+ if len(request.ConsumerServiceId) > 0 {
+ service, err = serviceUtil.GetService(ctx, domainProject,
request.ConsumerServiceId)
if err != nil {
log.Errorf(err, "get consumer failed, consumer[%s] find
provider instance[%s/%s]",
- in.ConsumerServiceId, in.ProviderServiceId,
in.ProviderInstanceId)
+ request.ConsumerServiceId,
request.ProviderServiceId, request.ProviderInstanceId)
return &pb.GetOneInstanceResponse{
Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
}, err
}
if service == nil {
log.Errorf(nil, "consumer does not exist, consumer[%s]
find provider instance[%s/%s]",
- in.ConsumerServiceId, in.ProviderServiceId,
in.ProviderInstanceId)
+ request.ConsumerServiceId,
request.ProviderServiceId, request.ProviderInstanceId)
return &pb.GetOneInstanceResponse{
Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
- fmt.Sprintf("Consumer[%s] does not
exist.", in.ConsumerServiceId)),
+ fmt.Sprintf("Consumer[%s] does not
exist.", request.ConsumerServiceId)),
}, nil
}
}
- provider, err := serviceUtil.GetService(ctx, domainProject,
in.ProviderServiceId)
+ provider, err := serviceUtil.GetService(ctx, domainProject,
request.ProviderServiceId)
if err != nil {
log.Errorf(err, "get provider failed, consumer[%s] find
provider instance[%s/%s]",
- in.ConsumerServiceId, in.ProviderServiceId,
in.ProviderInstanceId)
+ request.ConsumerServiceId, request.ProviderServiceId,
request.ProviderInstanceId)
return &pb.GetOneInstanceResponse{
Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
}, err
}
if provider == nil {
log.Errorf(nil, "provider does not exist, consumer[%s] find
provider instance[%s/%s]",
- in.ConsumerServiceId, in.ProviderServiceId,
in.ProviderInstanceId)
+ request.ConsumerServiceId, request.ProviderServiceId,
request.ProviderInstanceId)
return &pb.GetOneInstanceResponse{
Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
- fmt.Sprintf("Provider[%s] does not exist.",
in.ProviderServiceId)),
+ fmt.Sprintf("Provider[%s] does not exist.",
request.ProviderServiceId)),
}, nil
}
findFlag := func() string {
return fmt.Sprintf("Consumer[%s][%s/%s/%s/%s] find
provider[%s][%s/%s/%s/%s] instance[%s]",
- in.ConsumerServiceId, service.Environment,
service.AppId, service.ServiceName, service.Version,
+ request.ConsumerServiceId, service.Environment,
service.AppId, service.ServiceName, service.Version,
provider.ServiceId, provider.Environment,
provider.AppId, provider.ServiceName, provider.Version,
- in.ProviderInstanceId)
+ request.ProviderInstanceId)
}
var item *cache.VersionRuleCacheItem
rev, _ := ctx.Value(util.CtxRequestRevision).(string)
item, err = cache.FindInstances.GetWithProviderID(ctx, service,
proto.MicroServiceToKey(domainProject, provider),
&pb.HeartbeatSetElement{
- ServiceId: in.ProviderServiceId, InstanceId:
in.ProviderInstanceId,
- }, in.Tags, rev)
+ ServiceId: request.ProviderServiceId, InstanceId:
request.ProviderInstanceId,
+ }, request.Tags, rev)
if err != nil {
log.Errorf(err, "FindInstances.GetWithProviderID failed, %s
failed", findFlag())
return &pb.GetOneInstanceResponse{
@@ -433,20 +434,20 @@ func (ds *DataSource) GetInstance(ctx context.Context, in
*pb.GetOneInstanceRequ
}, nil
}
-func (ds *DataSource) FindInstances(ctx context.Context, in
*pb.FindInstancesRequest) (*pb.FindInstancesResponse,
+func (ds *DataSource) FindInstances(ctx context.Context, request
*pb.FindInstancesRequest) (*pb.FindInstancesResponse,
error) {
provider := &pb.MicroServiceKey{
Tenant: util.ParseTargetDomainProject(ctx),
- Environment: in.Environment,
- AppId: in.AppId,
- ServiceName: in.ServiceName,
- Alias: in.ServiceName,
- Version: in.VersionRule,
+ Environment: request.Environment,
+ AppId: request.AppId,
+ ServiceName: request.ServiceName,
+ Alias: request.ServiceName,
+ Version: request.VersionRule,
}
rev, ok := ctx.Value(util.CtxRequestRevision).(string)
if !ok {
- err := errors.New("rev in context is not type string")
+ err := errors.New("rev request context is not type string")
log.Error("", err)
return &pb.FindInstancesResponse{
Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
@@ -454,48 +455,48 @@ func (ds *DataSource) FindInstances(ctx context.Context,
in *pb.FindInstancesReq
}
if apt.IsShared(provider) {
- return ds.findSharedServiceInstance(ctx, in, provider, rev)
+ return ds.findSharedServiceInstance(ctx, request, provider, rev)
}
- return ds.findInstance(ctx, in, provider, rev)
+ return ds.findInstance(ctx, request, provider, rev)
}
-func (ds *DataSource) findInstance(ctx context.Context, in
*pb.FindInstancesRequest,
+func (ds *DataSource) findInstance(ctx context.Context, request
*pb.FindInstancesRequest,
provider *pb.MicroServiceKey, rev string) (*pb.FindInstancesResponse,
error) {
var err error
domainProject := util.ParseDomainProject(ctx)
- service := &pb.MicroService{Environment: in.Environment}
- if len(in.ConsumerServiceId) > 0 {
- service, err = serviceUtil.GetService(ctx, domainProject,
in.ConsumerServiceId)
+ service := &pb.MicroService{Environment: request.Environment}
+ if len(request.ConsumerServiceId) > 0 {
+ service, err = serviceUtil.GetService(ctx, domainProject,
request.ConsumerServiceId)
if err != nil {
log.Errorf(err, "get consumer failed, consumer[%s] find
provider[%s/%s/%s/%s]",
- in.ConsumerServiceId, in.Environment, in.AppId,
in.ServiceName, in.VersionRule)
+ request.ConsumerServiceId, request.Environment,
request.AppId, request.ServiceName, request.VersionRule)
return &pb.FindInstancesResponse{
Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
}, err
}
if service == nil {
log.Errorf(nil, "consumer does not exist, consumer[%s]
find provider[%s/%s/%s/%s]",
- in.ConsumerServiceId, in.Environment, in.AppId,
in.ServiceName, in.VersionRule)
+ request.ConsumerServiceId, request.Environment,
request.AppId, request.ServiceName, request.VersionRule)
return &pb.FindInstancesResponse{
Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
- fmt.Sprintf("Consumer[%s] does not
exist.", in.ConsumerServiceId)),
+ fmt.Sprintf("Consumer[%s] does not
exist.", request.ConsumerServiceId)),
}, nil
}
provider.Environment = service.Environment
}
// provider is not a shared micro-service,
- // only allow shared micro-service instances found in different domains.
+ // only allow shared micro-service instances found request different
domains.
ctx = util.SetTargetDomainProject(ctx, util.ParseDomain(ctx),
util.ParseProject(ctx))
provider.Tenant = util.ParseTargetDomainProject(ctx)
findFlag := fmt.Sprintf("Consumer[%s][%s/%s/%s/%s] find
provider[%s/%s/%s/%s]",
- in.ConsumerServiceId, service.Environment, service.AppId,
service.ServiceName, service.Version,
+ request.ConsumerServiceId, service.Environment, service.AppId,
service.ServiceName, service.Version,
provider.Environment, provider.AppId, provider.ServiceName,
provider.Version)
// cache
var item *cache.VersionRuleCacheItem
- item, err = cache.FindInstances.Get(ctx, service, provider, in.Tags,
rev)
+ item, err = cache.FindInstances.Get(ctx, service, provider,
request.Tags, rev)
if err != nil {
log.Errorf(err, "FindInstancesCache.Get failed, %s failed",
findFlag)
return &pb.FindInstancesResponse{
@@ -511,9 +512,9 @@ func (ds *DataSource) findInstance(ctx context.Context, in
*pb.FindInstancesRequ
}
// add dependency queue
- if len(in.ConsumerServiceId) > 0 &&
+ if len(request.ConsumerServiceId) > 0 &&
len(item.ServiceIds) > 0 &&
- !cache.DependencyRule.ExistVersionRule(ctx,
in.ConsumerServiceId, provider) {
+ !cache.DependencyRule.ExistVersionRule(ctx,
request.ConsumerServiceId, provider) {
provider, err = ds.reshapeProviderKey(ctx, provider,
item.ServiceIds[0])
if err != nil {
return nil, err
@@ -538,17 +539,17 @@ func (ds *DataSource) findInstance(ctx context.Context,
in *pb.FindInstancesRequ
return ds.genFindResult(ctx, rev, item)
}
-func (ds *DataSource) findSharedServiceInstance(ctx context.Context, in
*pb.FindInstancesRequest,
+func (ds *DataSource) findSharedServiceInstance(ctx context.Context, request
*pb.FindInstancesRequest,
provider *pb.MicroServiceKey, rev string) (*pb.FindInstancesResponse,
error) {
var err error
- service := &pb.MicroService{Environment: in.Environment}
+ service := &pb.MicroService{Environment: request.Environment}
// it means the shared micro-services must be the same env with SC.
provider.Environment = apt.Service.Environment
findFlag := fmt.Sprintf("find shared provider[%s/%s/%s/%s]",
provider.Environment, provider.AppId, provider.ServiceName, provider.Version)
// cache
var item *cache.VersionRuleCacheItem
- item, err = cache.FindInstances.Get(ctx, service, provider, in.Tags,
rev)
+ item, err = cache.FindInstances.Get(ctx, service, provider,
request.Tags, rev)
if err != nil {
log.Errorf(err, "FindInstancesCache.Get failed, %s failed",
findFlag)
return &pb.FindInstancesResponse{
@@ -594,12 +595,12 @@ func (ds *DataSource) reshapeProviderKey(ctx
context.Context, provider *pb.Micro
return provider, nil
}
-func (ds *DataSource) UpdateInstanceStatus(ctx context.Context, in
*pb.UpdateInstanceStatusRequest) (*pb.
+func (ds *DataSource) UpdateInstanceStatus(ctx context.Context, request
*pb.UpdateInstanceStatusRequest) (*pb.
UpdateInstanceStatusResponse, error) {
domainProject := util.ParseDomainProject(ctx)
- updateStatusFlag := util.StringJoin([]string{in.ServiceId,
in.InstanceId, in.Status}, "/")
+ updateStatusFlag := util.StringJoin([]string{request.ServiceId,
request.InstanceId, request.Status}, "/")
- instance, err := serviceUtil.GetInstance(ctx, domainProject,
in.ServiceId, in.InstanceId)
+ instance, err := serviceUtil.GetInstance(ctx, domainProject,
request.ServiceId, request.InstanceId)
if err != nil {
log.Errorf(err, "update instance[%s] status failed",
updateStatusFlag)
return &pb.UpdateInstanceStatusResponse{
@@ -614,7 +615,7 @@ func (ds *DataSource) UpdateInstanceStatus(ctx
context.Context, in *pb.UpdateIns
}
copyInstanceRef := *instance
- copyInstanceRef.Status = in.Status
+ copyInstanceRef.Status = request.Status
if err := serviceUtil.UpdateInstance(ctx, domainProject,
©InstanceRef); err != nil {
log.Errorf(err, "update instance[%s] status failed",
updateStatusFlag)
@@ -633,12 +634,12 @@ func (ds *DataSource) UpdateInstanceStatus(ctx
context.Context, in *pb.UpdateIns
}, nil
}
-func (ds *DataSource) UpdateInstanceProperties(ctx context.Context, in
*pb.UpdateInstancePropsRequest) (
+func (ds *DataSource) UpdateInstanceProperties(ctx context.Context, request
*pb.UpdateInstancePropsRequest) (
*pb.UpdateInstancePropsResponse, error) {
domainProject := util.ParseDomainProject(ctx)
- instanceFlag := util.StringJoin([]string{in.ServiceId, in.InstanceId},
"/")
+ instanceFlag := util.StringJoin([]string{request.ServiceId,
request.InstanceId}, "/")
- instance, err := serviceUtil.GetInstance(ctx, domainProject,
in.ServiceId, in.InstanceId)
+ instance, err := serviceUtil.GetInstance(ctx, domainProject,
request.ServiceId, request.InstanceId)
if err != nil {
log.Errorf(err, "update instance[%s] properties failed",
instanceFlag)
return &pb.UpdateInstancePropsResponse{
@@ -653,7 +654,7 @@ func (ds *DataSource) UpdateInstanceProperties(ctx
context.Context, in *pb.Updat
}
copyInstanceRef := *instance
- copyInstanceRef.Properties = in.Properties
+ copyInstanceRef.Properties = request.Properties
if err := serviceUtil.UpdateInstance(ctx, domainProject,
©InstanceRef); err != nil {
log.Errorf(err, "update instance[%s] properties failed",
instanceFlag)
@@ -672,16 +673,16 @@ func (ds *DataSource) UpdateInstanceProperties(ctx
context.Context, in *pb.Updat
}, nil
}
-func (ds *DataSource) HeartbeatSet(ctx context.Context, in
*pb.HeartbeatSetRequest) (*pb.HeartbeatSetResponse, error) {
+func (ds *DataSource) HeartbeatSet(ctx context.Context, request
*pb.HeartbeatSetRequest) (*pb.HeartbeatSetResponse, error) {
domainProject := util.ParseDomainProject(ctx)
- heartBeatCount := len(in.Instances)
+ heartBeatCount := len(request.Instances)
existFlag := make(map[string]bool, heartBeatCount)
instancesHbRst := make(chan *pb.InstanceHbRst, heartBeatCount)
noMultiCounter := 0
- for _, heartbeatElement := range in.Instances {
+ for _, heartbeatElement := range request.Instances {
if _, ok :=
existFlag[heartbeatElement.ServiceId+heartbeatElement.InstanceId]; ok {
- log.Warnf("instance[%s/%s] is duplicate in heartbeat
set",
+ log.Warnf("instance[%s/%s] is duplicate request
heartbeat set",
heartbeatElement.ServiceId,
heartbeatElement.InstanceId)
continue
} else {
@@ -713,58 +714,58 @@ func (ds *DataSource) HeartbeatSet(ctx context.Context,
in *pb.HeartbeatSetReque
Instances: instanceHbRstArr,
}, nil
}
- log.Errorf(nil, "batch update heartbeats failed, %v", in.Instances)
+ log.Errorf(nil, "batch update heartbeats failed, %v", request.Instances)
return &pb.HeartbeatSetResponse{
Response: proto.CreateResponse(scerr.ErrInstanceNotExists,
"Heartbeat set failed."),
Instances: instanceHbRstArr,
}, nil
}
-func (ds *DataSource) GetInstances(ctx context.Context, in
*pb.GetInstancesRequest) (*pb.GetInstancesResponse,
+func (ds *DataSource) GetInstances(ctx context.Context, request
*pb.GetInstancesRequest) (*pb.GetInstancesResponse,
error) {
domainProject := util.ParseDomainProject(ctx)
service := &pb.MicroService{}
var err error
- if len(in.ConsumerServiceId) > 0 {
- service, err = serviceUtil.GetService(ctx, domainProject,
in.ConsumerServiceId)
+ if len(request.ConsumerServiceId) > 0 {
+ service, err = serviceUtil.GetService(ctx, domainProject,
request.ConsumerServiceId)
if err != nil {
log.Errorf(err, "get consumer failed, consumer[%s] find
provider instances",
- in.ConsumerServiceId, in.ProviderServiceId)
+ request.ConsumerServiceId,
request.ProviderServiceId)
return &pb.GetInstancesResponse{
Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
}, err
}
if service == nil {
log.Errorf(nil, "consumer does not exist, consumer[%s]
find provider instances",
- in.ConsumerServiceId, in.ProviderServiceId)
+ request.ConsumerServiceId,
request.ProviderServiceId)
return &pb.GetInstancesResponse{
Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
- fmt.Sprintf("Consumer[%s] does not
exist.", in.ConsumerServiceId)),
+ fmt.Sprintf("Consumer[%s] does not
exist.", request.ConsumerServiceId)),
}, nil
}
}
- provider, err := serviceUtil.GetService(ctx, domainProject,
in.ProviderServiceId)
+ provider, err := serviceUtil.GetService(ctx, domainProject,
request.ProviderServiceId)
if err != nil {
log.Errorf(err, "get provider failed, consumer[%s] find
provider instances",
- in.ConsumerServiceId, in.ProviderServiceId)
+ request.ConsumerServiceId, request.ProviderServiceId)
return &pb.GetInstancesResponse{
Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
}, err
}
if provider == nil {
log.Errorf(nil, "provider does not exist, consumer[%s] find
provider instances",
- in.ConsumerServiceId, in.ProviderServiceId)
+ request.ConsumerServiceId, request.ProviderServiceId)
return &pb.GetInstancesResponse{
Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
- fmt.Sprintf("Provider[%s] does not exist.",
in.ProviderServiceId)),
+ fmt.Sprintf("Provider[%s] does not exist.",
request.ProviderServiceId)),
}, nil
}
findFlag := func() string {
return fmt.Sprintf("Consumer[%s][%s/%s/%s/%s] find
provider[%s][%s/%s/%s/%s] instances",
- in.ConsumerServiceId, service.Environment,
service.AppId, service.ServiceName, service.Version,
+ request.ConsumerServiceId, service.Environment,
service.AppId, service.ServiceName, service.Version,
provider.ServiceId, provider.Environment,
provider.AppId, provider.ServiceName, provider.Version)
}
@@ -772,8 +773,8 @@ func (ds *DataSource) GetInstances(ctx context.Context, in
*pb.GetInstancesReque
rev, _ := ctx.Value(util.CtxRequestRevision).(string)
item, err = cache.FindInstances.GetWithProviderID(ctx, service,
proto.MicroServiceToKey(domainProject, provider),
&pb.HeartbeatSetElement{
- ServiceId: in.ProviderServiceId,
- }, in.Tags, rev)
+ ServiceId: request.ProviderServiceId,
+ }, request.Tags, rev)
if err != nil {
log.Errorf(err, "FindInstances.GetWithProviderID failed, %s
failed", findFlag())
return &pb.GetInstancesResponse{
@@ -826,19 +827,19 @@ func (ds *DataSource) BatchFind(ctx context.Context,
request *pb.BatchFindInstan
return response, nil
}
-func (ds *DataSource) batchFindServices(ctx context.Context, in
*pb.BatchFindInstancesRequest) (
+func (ds *DataSource) batchFindServices(ctx context.Context, request
*pb.BatchFindInstancesRequest) (
*pb.BatchFindResult, error) {
- if len(in.Services) == 0 {
+ if len(request.Services) == 0 {
return nil, nil
}
cloneCtx := util.CloneContext(ctx)
services := &pb.BatchFindResult{}
failedResult := make(map[int32]*pb.FindFailedResult)
- for index, key := range in.Services {
+ for index, key := range request.Services {
findCtx := util.SetContext(cloneCtx, util.CtxRequestRevision,
key.Rev)
resp, err := ds.FindInstances(findCtx, &pb.FindInstancesRequest{
- ConsumerServiceId: in.ConsumerServiceId,
+ ConsumerServiceId: request.ConsumerServiceId,
AppId: key.Service.AppId,
ServiceName: key.Service.ServiceName,
VersionRule: key.Service.Version,
@@ -860,8 +861,8 @@ func (ds *DataSource) batchFindServices(ctx
context.Context, in *pb.BatchFindIns
return services, nil
}
-func (ds *DataSource) batchFindInstances(ctx context.Context, in
*pb.BatchFindInstancesRequest) (*pb.BatchFindResult, error) {
- if len(in.Instances) == 0 {
+func (ds *DataSource) batchFindInstances(ctx context.Context, request
*pb.BatchFindInstancesRequest) (*pb.BatchFindResult, error) {
+ if len(request.Instances) == 0 {
return nil, nil
}
cloneCtx := util.CloneContext(ctx)
@@ -870,10 +871,10 @@ func (ds *DataSource) batchFindInstances(ctx
context.Context, in *pb.BatchFindIn
instances := &pb.BatchFindResult{}
failedResult := make(map[int32]*pb.FindFailedResult)
- for index, key := range in.Instances {
+ for index, key := range request.Instances {
getCtx := util.SetContext(cloneCtx, util.CtxRequestRevision,
key.Rev)
resp, err := ds.GetInstance(getCtx, &pb.GetOneInstanceRequest{
- ConsumerServiceId: in.ConsumerServiceId,
+ ConsumerServiceId: request.ConsumerServiceId,
ProviderServiceId: key.Instance.ServiceId,
ProviderInstanceId: key.Instance.InstanceId,
})
@@ -1020,14 +1021,6 @@ func (ds *DataSource) ModifySchema(ctx context.Context,
request *pb.ModifySchema
}, nil
}
-func (ds *DataSource) GetSchema() {
- panic("implement me")
-}
-
-func (ds *DataSource) DeleteSchema() {
- panic("implement me")
-}
-
func (ds *DataSource) ExistSchema(ctx context.Context, request
*pb.GetExistenceRequest) (
*pb.GetExistenceResponse, error) {
domainProject := util.ParseDomainProject(ctx)
@@ -1068,25 +1061,206 @@ func (ds *DataSource) ExistSchema(ctx context.Context,
request *pb.GetExistenceR
}, nil
}
-func (ds *DataSource) AddTags(ctx context.Context, in
*pb.AddServiceTagsRequest) (*pb.AddServiceTagsResponse, error) {
+func (ds *DataSource) GetSchema(ctx context.Context, request
*pb.GetSchemaRequest) (*pb.GetSchemaResponse, error) {
+ domainProject := util.ParseDomainProject(ctx)
+
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(nil, "get schema[%s/%s] failed, service does not
exist",
+ request.ServiceId, request.SchemaId)
+ return &pb.GetSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ key := apt.GenerateServiceSchemaKey(domainProject, request.ServiceId,
request.SchemaId)
+ opts := append(serviceUtil.FromContext(ctx), registry.WithStrKey(key))
+ resp, errDo := backend.Store().Schema().Search(ctx, opts...)
+ if errDo != nil {
+ log.Errorf(errDo, "get schema[%s/%s] failed",
request.ServiceId, request.SchemaId)
+ return &pb.GetSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, errDo.Error()),
+ }, errDo
+ }
+ if resp.Count == 0 {
+ log.Errorf(errDo, "get schema[%s/%s] failed, schema does not
exists",
+ request.ServiceId, request.SchemaId)
+ return &pb.GetSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrSchemaNotExists, "Do not have this schema info."),
+ }, nil
+ }
+
+ schemaSummary, err := getSchemaSummary(ctx, domainProject,
request.ServiceId, request.SchemaId)
+ if err != nil {
+ log.Errorf(err, "get schema[%s/%s] failed, get schema summary
failed",
+ request.ServiceId, request.SchemaId)
+ return &pb.GetSchemaResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ return &pb.GetSchemaResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
"Get schema info successfully."),
+ Schema:
util.BytesToStringWithNoCopy(resp.Kvs[0].Value.([]byte)),
+ SchemaSummary: schemaSummary,
+ }, nil
+}
+
+func (ds *DataSource) GetAllSchemas(ctx context.Context, request
*pb.GetAllSchemaRequest) (
+ *pb.GetAllSchemaResponse, error) {
+ domainProject := util.ParseDomainProject(ctx)
+
+ service, err := serviceUtil.GetService(ctx, domainProject,
request.ServiceId)
+ if err != nil {
+ log.Errorf(err, "get service[%s] all schemas failed, get
service failed", request.ServiceId)
+ return &pb.GetAllSchemaResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if service == nil {
+ log.Errorf(nil, "get service[%s] all schemas failed, service
does not exist", request.ServiceId)
+ return &pb.GetAllSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ schemasList := service.Schemas
+ if len(schemasList) == 0 {
+ return &pb.GetAllSchemaResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
"Do not have this schema info."),
+ Schemas: []*pb.Schema{},
+ }, nil
+ }
+
+ key := apt.GenerateServiceSchemaSummaryKey(domainProject,
request.ServiceId, "")
+ opts := append(serviceUtil.FromContext(ctx), registry.WithStrKey(key),
registry.WithPrefix())
+ resp, errDo := backend.Store().SchemaSummary().Search(ctx, opts...)
+ if errDo != nil {
+ log.Errorf(errDo, "get service[%s] all schema summaries
failed", request.ServiceId)
+ return &pb.GetAllSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, errDo.Error()),
+ }, errDo
+ }
+
+ respWithSchema := &discovery.Response{}
+ if request.WithSchema {
+ key := apt.GenerateServiceSchemaKey(domainProject,
request.ServiceId, "")
+ opts := append(serviceUtil.FromContext(ctx),
registry.WithStrKey(key), registry.WithPrefix())
+ respWithSchema, errDo = backend.Store().Schema().Search(ctx,
opts...)
+ if errDo != nil {
+ log.Errorf(errDo, "get service[%s] all schemas failed",
request.ServiceId)
+ return &pb.GetAllSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, errDo.Error()),
+ }, errDo
+ }
+ }
+
+ schemas := make([]*pb.Schema, 0, len(schemasList))
+ for _, schemaID := range schemasList {
+ tempSchema := &pb.Schema{}
+ tempSchema.SchemaId = schemaID
+ for _, summarySchema := range resp.Kvs {
+ _, _, schemaIDOfSummary :=
apt.GetInfoFromSchemaSummaryKV(summarySchema.Key)
+ if schemaID == schemaIDOfSummary {
+ tempSchema.Summary =
summarySchema.Value.(string)
+ }
+ }
+
+ for _, contentSchema := range respWithSchema.Kvs {
+ _, _, schemaIDOfSchema :=
apt.GetInfoFromSchemaKV(contentSchema.Key)
+ if schemaID == schemaIDOfSchema {
+ tempSchema.Schema =
util.BytesToStringWithNoCopy(contentSchema.Value.([]byte))
+ }
+ }
+ schemas = append(schemas, tempSchema)
+ }
+
+ return &pb.GetAllSchemaResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Get all
schema info successfully."),
+ Schemas: schemas,
+ }, nil
+}
+
+func (ds *DataSource) DeleteSchema(ctx context.Context, request
*pb.DeleteSchemaRequest) (
+ *pb.DeleteSchemaResponse, error) {
+ remoteIP := util.GetIPFromContext(ctx)
+ domainProject := util.ParseDomainProject(ctx)
+
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(nil, "delete schema[%s/%s] failed, service does not
exist, operator: %s",
+ request.ServiceId, request.SchemaId, remoteIP)
+ return &pb.DeleteSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ key := apt.GenerateServiceSchemaKey(domainProject, request.ServiceId,
request.SchemaId)
+ exist, err := serviceUtil.CheckSchemaInfoExist(ctx, key)
+ if err != nil {
+ log.Errorf(err, "delete schema[%s/%s] failed, operator: %s",
+ request.ServiceId, request.SchemaId, remoteIP)
+ return &pb.DeleteSchemaResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if !exist {
+ log.Errorf(nil, "delete schema[%s/%s] failed, schema does not
exist, operator: %s",
+ request.ServiceId, request.SchemaId, remoteIP)
+ return &pb.DeleteSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrSchemaNotExists, "Schema info does not exist."),
+ }, nil
+ }
+ epSummaryKey := apt.GenerateServiceSchemaSummaryKey(domainProject,
request.ServiceId, request.SchemaId)
+ opts := []registry.PluginOp{
+ registry.OpDel(registry.WithStrKey(epSummaryKey)),
+ registry.OpDel(registry.WithStrKey(key)),
+ }
+
+ resp, errDo := backend.Registry().TxnWithCmp(ctx, opts,
+ []registry.CompareOp{registry.OpCmp(
+
registry.CmpVer(util.StringToBytesWithNoCopy(apt.GenerateServiceKey(domainProject,
request.ServiceId))),
+ registry.CmpNotEqual, 0)},
+ nil)
+ if errDo != nil {
+ log.Errorf(errDo, "delete schema[%s/%s] failed, operator: %s",
+ request.ServiceId, request.SchemaId, remoteIP)
+ return &pb.DeleteSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, errDo.Error()),
+ }, errDo
+ }
+ if !resp.Succeeded {
+ log.Errorf(nil, "delete schema[%s/%s] failed, service does not
exist, operator: %s",
+ request.ServiceId, request.SchemaId, remoteIP)
+ return &pb.DeleteSchemaResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ log.Infof("delete schema[%s/%s] info successfully, operator: %s",
+ request.ServiceId, request.SchemaId, remoteIP)
+ return &pb.DeleteSchemaResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Delete
schema info successfully."),
+ }, nil
+}
+
+func (ds *DataSource) AddTags(ctx context.Context, request
*pb.AddServiceTagsRequest) (*pb.AddServiceTagsResponse, error) {
remoteIP := util.GetIPFromContext(ctx)
domainProject := util.ParseDomainProject(ctx)
// service id存在性校验
- if !serviceUtil.ServiceExist(ctx, domainProject, in.ServiceId) {
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
log.Errorf(nil, "add service[%s]'s tags %v failed, service does
not exist, operator: %s",
- in.ServiceId, in.Tags, remoteIP)
+ request.ServiceId, request.Tags, remoteIP)
return &pb.AddServiceTagsResponse{
Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
}, nil
}
- addTags := in.Tags
- res := quota.NewApplyQuotaResource(quota.TagQuotaType, domainProject,
in.ServiceId, int64(len(addTags)))
+ addTags := request.Tags
+ res := quota.NewApplyQuotaResource(quota.TagQuotaType, domainProject,
request.ServiceId, int64(len(addTags)))
rst := plugin.Plugins().Quota().Apply4Quotas(ctx, res)
errQuota := rst.Err
if errQuota != nil {
- log.Errorf(errQuota, "add service[%s]'s tags %v failed,
operator: %s", in.ServiceId, addTags, remoteIP)
+ log.Errorf(errQuota, "add service[%s]'s tags %v failed,
operator: %s", request.ServiceId, addTags, remoteIP)
response := &pb.AddServiceTagsResponse{
Response: proto.CreateResponseWithSCErr(errQuota),
}
@@ -1096,10 +1270,10 @@ func (ds *DataSource) AddTags(ctx context.Context, in
*pb.AddServiceTagsRequest)
return response, nil
}
- dataTags, err := serviceUtil.GetTagsUtils(ctx, domainProject,
in.ServiceId)
+ dataTags, err := serviceUtil.GetTagsUtils(ctx, domainProject,
request.ServiceId)
if err != nil {
log.Errorf(err, "add service[%s]'s tags %v failed, get existed
tag failed, operator: %s",
- in.ServiceId, addTags, remoteIP)
+ request.ServiceId, addTags, remoteIP)
return &pb.AddServiceTagsResponse{
Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
}, err
@@ -1112,9 +1286,9 @@ func (ds *DataSource) AddTags(ctx context.Context, in
*pb.AddServiceTagsRequest)
}
dataTags = addTags
- checkErr := serviceUtil.AddTagIntoETCD(ctx, domainProject,
in.ServiceId, dataTags)
+ checkErr := serviceUtil.AddTagIntoETCD(ctx, domainProject,
request.ServiceId, dataTags)
if checkErr != nil {
- log.Errorf(checkErr, "add service[%s]'s tags %v failed,
operator: %s", in.ServiceId, in.Tags, remoteIP)
+ log.Errorf(checkErr, "add service[%s]'s tags %v failed,
operator: %s", request.ServiceId, request.Tags, remoteIP)
resp := &pb.AddServiceTagsResponse{
Response: proto.CreateResponseWithSCErr(checkErr),
}
@@ -1124,7 +1298,7 @@ func (ds *DataSource) AddTags(ctx context.Context, in
*pb.AddServiceTagsRequest)
return resp, nil
}
- log.Infof("add service[%s]'s tags %v successfully, operator: %s",
in.ServiceId, in.Tags, remoteIP)
+ log.Infof("add service[%s]'s tags %v successfully, operator: %s",
request.ServiceId, request.Tags, remoteIP)
return &pb.AddServiceTagsResponse{
Response: proto.CreateResponse(proto.Response_SUCCESS, "Add
service tags successfully."),
}, nil
@@ -1142,20 +1316,331 @@ func (ds *DataSource) DeleteTag() {
panic("implement me")
}
-func (ds *DataSource) AddRule() {
- panic("implement me")
+func (ds *DataSource) AddRule(ctx context.Context, request
*pb.AddServiceRulesRequest) (
+ *pb.AddServiceRulesResponse, error) {
+ remoteIP := util.GetIPFromContext(ctx)
+ domainProject := util.ParseDomainProject(ctx)
+
+ // service id存在性校验
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(nil, "add service[%s] rule failed, service does not
exist, operator: %s",
+ request.ServiceId, remoteIP)
+ return &pb.AddServiceRulesResponse{
+ Response: proto.CreateResponse(scerr.ErrInvalidParams,
"Service does not exist."),
+ }, nil
+ }
+ res := quota.NewApplyQuotaResource(quota.RuleQuotaType, domainProject,
request.ServiceId, int64(len(request.Rules)))
+ rst := plugin.Plugins().Quota().Apply4Quotas(ctx, res)
+ errQuota := rst.Err
+ if errQuota != nil {
+ log.Errorf(errQuota, "add service[%s] rule failed, operator:
%s", request.ServiceId, remoteIP)
+ response := &pb.AddServiceRulesResponse{
+ Response: proto.CreateResponseWithSCErr(errQuota),
+ }
+ if errQuota.InternalError() {
+ return response, errQuota
+ }
+ return response, nil
+ }
+
+ ruleType, _, err := serviceUtil.GetServiceRuleType(ctx, domainProject,
request.ServiceId)
+ if err != nil {
+ return &pb.AddServiceRulesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ ruleIDs := make([]string, 0, len(request.Rules))
+ opts := make([]registry.PluginOp, 0, 2*len(request.Rules))
+ for _, rule := range request.Rules {
+ //黑白名单只能存在一种,黑名单 or 白名单
+ if len(ruleType) == 0 {
+ ruleType = rule.RuleType
+ } else if ruleType != rule.RuleType {
+ log.Errorf(nil,
+ "add service[%s] rule failed, can not add
different RuleType at the same time, operator: %s",
+ request.ServiceId, remoteIP)
+ return &pb.AddServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrBlackAndWhiteRule,
+ "Service can only contain one rule
type, BLACK or WHITE."),
+ }, nil
+
+ }
+
+ //同一服务,attribute和pattern确定一个rule
+ if serviceUtil.RuleExist(ctx, domainProject, request.ServiceId,
rule.Attribute, rule.Pattern) {
+ log.Infof("service[%s] rule[%s/%s] already exists,
operator: %s",
+ request.ServiceId, rule.Attribute,
rule.Pattern, remoteIP)
+ continue
+ }
+
+ // 产生全局rule id
+ timestamp := strconv.FormatInt(time.Now().Unix(), 10)
+ ruleAdd := &pb.ServiceRule{
+ RuleId: util.GenerateUUID(),
+ RuleType: rule.RuleType,
+ Attribute: rule.Attribute,
+ Pattern: rule.Pattern,
+ Description: rule.Description,
+ Timestamp: timestamp,
+ ModTimestamp: timestamp,
+ }
+
+ key := apt.GenerateServiceRuleKey(domainProject,
request.ServiceId, ruleAdd.RuleId)
+ indexKey := apt.GenerateRuleIndexKey(domainProject,
request.ServiceId, ruleAdd.Attribute, ruleAdd.Pattern)
+ ruleIDs = append(ruleIDs, ruleAdd.RuleId)
+
+ data, err := json.Marshal(ruleAdd)
+ if err != nil {
+ log.Errorf(err, "add service[%s] rule failed, marshal
rule[%s/%s] failed, operator: %s",
+ request.ServiceId, ruleAdd.Attribute,
ruleAdd.Pattern, remoteIP)
+ return &pb.AddServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
+ }, err
+ }
+
+ opts = append(opts, registry.OpPut(registry.WithStrKey(key),
registry.WithValue(data)))
+ opts = append(opts,
registry.OpPut(registry.WithStrKey(indexKey),
registry.WithStrValue(ruleAdd.RuleId)))
+ }
+ if len(opts) <= 0 {
+ log.Infof("add service[%s] rule successfully, no rules to add,
operator: %s",
+ request.ServiceId, remoteIP)
+ return &pb.AddServiceRulesResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
"Service rules has been added."),
+ }, nil
+ }
+
+ resp, err := backend.BatchCommitWithCmp(ctx, opts,
+ []registry.CompareOp{registry.OpCmp(
+
registry.CmpVer(util.StringToBytesWithNoCopy(apt.GenerateServiceKey(domainProject,
request.ServiceId))),
+ registry.CmpNotEqual, 0)},
+ nil)
+ if err != nil {
+ log.Errorf(err, "add service[%s] rule failed, operator: %s",
request.ServiceId, remoteIP)
+ return &pb.AddServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, err.Error()),
+ }, err
+ }
+ if !resp.Succeeded {
+ log.Errorf(nil, "add service[%s] rule failed, service does not
exist, operator: %s",
+ request.ServiceId, remoteIP)
+ return &pb.AddServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ log.Infof("add service[%s] rule %v successfully, operator: %s",
request.ServiceId, ruleIDs, remoteIP)
+ return &pb.AddServiceRulesResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Add
service rules successfully."),
+ RuleIds: ruleIDs,
+ }, nil
}
-func (ds *DataSource) GetRule() {
- panic("implement me")
+func (ds *DataSource) GetRule(ctx context.Context, request
*pb.GetServiceRulesRequest) (
+ *pb.GetServiceRulesResponse, error) {
+ domainProject := util.ParseDomainProject(ctx)
+
+ // service id存在性校验
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(nil, "get service[%s] rule failed, service does not
exist", request.ServiceId)
+ return &pb.GetServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ rules, err := serviceUtil.GetRulesUtil(ctx, domainProject,
request.ServiceId)
+ if err != nil {
+ log.Errorf(err, "get service[%s] rule failed",
request.ServiceId)
+ return &pb.GetServiceRulesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ return &pb.GetServiceRulesResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Get
service rules successfully."),
+ Rules: rules,
+ }, nil
}
-func (ds *DataSource) UpdateRule() {
- panic("implement me")
+func (ds *DataSource) UpdateRule(ctx context.Context, request
*pb.UpdateServiceRuleRequest) (
+ *pb.UpdateServiceRuleResponse, error) {
+ remoteIP := util.GetIPFromContext(ctx)
+ domainProject := util.ParseDomainProject(ctx)
+
+ // service id存在性校验
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(nil, "update service rule[%s/%s] failed, service
does not exist, operator: %s",
+ request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ //是否能改变ruleType
+ ruleType, ruleNum, err := serviceUtil.GetServiceRuleType(ctx,
domainProject, request.ServiceId)
+ if err != nil {
+ log.Errorf(err, "update service rule[%s/%s] failed, get rule
type failed, operator: %s",
+ request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if ruleNum >= 1 && ruleType != request.Rule.RuleType {
+ log.Errorf(err, "update service rule[%s/%s] failed, can only
exist one type, current type is %s, operator: %s",
+ request.ServiceId, request.RuleId, ruleType, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response:
proto.CreateResponse(scerr.ErrModifyRuleNotAllow, "Exist multiple rules,can not
change rule type. Rule type is "+ruleType),
+ }, nil
+ }
+
+ rule, err := serviceUtil.GetOneRule(ctx, domainProject,
request.ServiceId, request.RuleId)
+ if err != nil {
+ log.Errorf(err, "update service rule[%s/%s] failed, query
service rule failed, operator: %s",
+ request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if rule == nil {
+ log.Errorf(err, "update service rule[%s/%s] failed, service
rule does not exist, operator: %s",
+ request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response: proto.CreateResponse(scerr.ErrRuleNotExists,
"This rule does not exist."),
+ }, nil
+ }
+
+ copyRuleRef := *rule
+ oldRulePatten := copyRuleRef.Pattern
+ oldRuleAttr := copyRuleRef.Attribute
+ isChangeIndex := false
+ if copyRuleRef.Attribute != request.Rule.Attribute {
+ isChangeIndex = true
+ copyRuleRef.Attribute = request.Rule.Attribute
+ }
+ if copyRuleRef.Pattern != request.Rule.Pattern {
+ isChangeIndex = true
+ copyRuleRef.Pattern = request.Rule.Pattern
+ }
+ copyRuleRef.RuleType = request.Rule.RuleType
+ copyRuleRef.Description = request.Rule.Description
+ copyRuleRef.ModTimestamp = strconv.FormatInt(time.Now().Unix(), 10)
+
+ key := apt.GenerateServiceRuleKey(domainProject, request.ServiceId,
request.RuleId)
+ data, err := json.Marshal(copyRuleRef)
+ if err != nil {
+ log.Errorf(err, "update service rule[%s/%s] failed, marshal
service rule failed, operator: %s",
+ request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ var opts []registry.PluginOp
+ if isChangeIndex {
+ //加入新的rule index
+ indexKey := apt.GenerateRuleIndexKey(domainProject,
request.ServiceId, copyRuleRef.Attribute, copyRuleRef.Pattern)
+ opts = append(opts,
registry.OpPut(registry.WithStrKey(indexKey),
registry.WithStrValue(copyRuleRef.RuleId)))
+
+ //删除旧的rule index
+ oldIndexKey := apt.GenerateRuleIndexKey(domainProject,
request.ServiceId, oldRuleAttr, oldRulePatten)
+ opts = append(opts,
registry.OpDel(registry.WithStrKey(oldIndexKey)))
+ }
+ opts = append(opts, registry.OpPut(registry.WithStrKey(key),
registry.WithValue(data)))
+
+ resp, err := backend.Registry().TxnWithCmp(ctx, opts,
+ []registry.CompareOp{registry.OpCmp(
+
registry.CmpVer(util.StringToBytesWithNoCopy(apt.GenerateServiceKey(domainProject,
request.ServiceId))),
+ registry.CmpNotEqual, 0)},
+ nil)
+ if err != nil {
+ log.Errorf(err, "update service rule[%s/%s] failed, operator:
%s", request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, err.Error()),
+ }, err
+ }
+ if !resp.Succeeded {
+ log.Errorf(err, "update service rule[%s/%s] failed, service
does not exist, operator: %s",
+ request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ log.Infof("update service rule[%s/%s] successfully, operator: %s",
request.ServiceId, request.RuleId, remoteIP)
+ return &pb.UpdateServiceRuleResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Get
service rules successfully."),
+ }, nil
}
-func (ds *DataSource) DeleteRule() {
- panic("implement me")
+func (ds *DataSource) DeleteRule(ctx context.Context, request
*pb.DeleteServiceRulesRequest) (
+ *pb.DeleteServiceRulesResponse, error) {
+ remoteIP := util.GetIPFromContext(ctx)
+ domainProject := util.ParseDomainProject(ctx)
+
+ // service id存在性校验
+ if !serviceUtil.ServiceExist(ctx, domainProject, request.ServiceId) {
+ log.Errorf(nil, "delete service[%s] rules %v failed, service
does not exist, operator: %s",
+ request.ServiceId, request.RuleIds, remoteIP)
+ return &pb.DeleteServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ opts := []registry.PluginOp{}
+ key := ""
+ indexKey := ""
+ for _, ruleID := range request.RuleIds {
+ key = apt.GenerateServiceRuleKey(domainProject,
request.ServiceId, ruleID)
+ log.Debugf("start delete service rule file: %s", key)
+ data, err := serviceUtil.GetOneRule(ctx, domainProject,
request.ServiceId, ruleID)
+ if err != nil {
+ log.Errorf(err, "delete service[%s] rules %v failed,
get rule[%s] failed, operator: %s",
+ request.ServiceId, request.RuleIds, ruleID,
remoteIP)
+ return &pb.DeleteServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
+ }, err
+ }
+ if data == nil {
+ log.Errorf(nil, "delete service[%s] rules %v failed,
rule[%s] does not exist, operator: %s",
+ request.ServiceId, request.RuleIds, ruleID,
remoteIP)
+ return &pb.DeleteServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrRuleNotExists, "This rule does not exist."),
+ }, nil
+ }
+ indexKey = apt.GenerateRuleIndexKey(domainProject,
request.ServiceId, data.Attribute, data.Pattern)
+ opts = append(opts,
+ registry.OpDel(registry.WithStrKey(key)),
+ registry.OpDel(registry.WithStrKey(indexKey)))
+ }
+ if len(opts) <= 0 {
+ log.Errorf(nil, "delete service[%s] rules %v failed, no rule
has been deleted, operator: %s",
+ request.ServiceId, request.RuleIds, remoteIP)
+ return &pb.DeleteServiceRulesResponse{
+ Response: proto.CreateResponse(scerr.ErrRuleNotExists,
"No service rule has been deleted."),
+ }, nil
+ }
+
+ resp, err := backend.BatchCommitWithCmp(ctx, opts,
+ []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] rules %v failed, operator:
%s", request.ServiceId, request.RuleIds, remoteIP)
+ return &pb.DeleteServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrUnavailableBackend, err.Error()),
+ }, err
+ }
+ if !resp.Succeeded {
+ log.Errorf(err, "delete service[%s] rules %v failed, service
does not exist, operator: %s",
+ request.ServiceId, request.RuleIds, remoteIP)
+ return &pb.DeleteServiceRulesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, "Service does not exist."),
+ }, nil
+ }
+
+ log.Infof("delete service[%s] rules %v successfully, operator: %s",
request.ServiceId, request.RuleIds, remoteIP)
+ return &pb.DeleteServiceRulesResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Delete
service rules successfully."),
+ }, nil
}
func (ds *DataSource) modifySchemas(ctx context.Context, domainProject string,
service *pb.MicroService,
@@ -1302,7 +1787,7 @@ func (ds *DataSource) modifySchema(ctx context.Context,
serviceID string, schema
if !ds.isSchemaEditable(microService) {
if len(microService.Schemas) != 0 && !isExist {
- return scerr.NewError(scerr.ErrUndefinedSchemaID,
"Non-existent schemaID can't be added in "+pb.ENV_PROD)
+ return scerr.NewError(scerr.ErrUndefinedSchemaID,
"Non-existent schemaID can't be added request "+pb.ENV_PROD)
}
key := apt.GenerateServiceSchemaKey(domainProject, serviceID,
schemaID)
@@ -1318,7 +1803,7 @@ func (ds *DataSource) modifySchema(ctx context.Context,
serviceID string, schema
log.Errorf(err, "%s mode, schema[%s/%s] already
exists, can not be changed, operator: %s",
pb.ENV_PROD, serviceID, schemaID,
remoteIP)
return
scerr.NewError(scerr.ErrModifySchemaNotAllow,
- "schema already exist, can not be
changed in "+pb.ENV_PROD)
+ "schema already exist, can not be
changed request "+pb.ENV_PROD)
}
exist, err := isExistSchemaSummary(ctx, domainProject,
serviceID, schemaID)
@@ -1330,7 +1815,7 @@ func (ds *DataSource) modifySchema(ctx context.Context,
serviceID string, schema
if exist {
log.Errorf(err, "%s mode, schema[%s/%s] already
exist, can not be changed, operator: %s",
pb.ENV_PROD, serviceID, schemaID,
remoteIP)
- return
scerr.NewError(scerr.ErrModifySchemaNotAllow, "schema already exist, can not be
changed in "+pb.ENV_PROD)
+ return
scerr.NewError(scerr.ErrModifySchemaNotAllow, "schema already exist, can not be
changed request "+pb.ENV_PROD)
}
}
diff --git a/datasource/etcd/ms_test.go b/datasource/etcd/ms_test.go
index f358f8c..948be0a 100644
--- a/datasource/etcd/ms_test.go
+++ b/datasource/etcd/ms_test.go
@@ -2776,11 +2776,825 @@ func TestSchema_Exist(t *testing.T) {
})
}
-func TestTag_Add(t *testing.T) {
+func TestSchema_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)
+
+ var (
+ serviceId string
+ serviceId1 string
+ )
+
+ var (
+ schemaId1 string = "all_schema1_ms"
+ schemaId2 string = "all_schema2_ms"
+ schemaId3 string = "all_schema3_ms"
+ summary string = "this0is1a2test3ms"
+ schemaContent string = "the content is vary large"
+ )
+
+ t.Run("register service and instance", func(t *testing.T) {
+ respCreateService, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "get_schema_group_ms",
+ ServiceName: "get_schema_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Schemas: []string{
+ "non-schema-content",
+ },
+ Status: pb.MS_UP,
+ Environment: pb.ENV_DEV,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId = respCreateService.ServiceId
+
+ respCreateSchema, err :=
datasource.Instance().ModifySchema(getContext(), &pb.ModifySchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "com.huawei.test.ms",
+ Schema: "get schema ms",
+ Summary: "schema0summary1ms",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateSchema.Response.GetCode())
+
+ respCreateService, err =
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "get_all_schema_ms",
+ ServiceName: "get_all_schema_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Schemas: []string{
+ schemaId1,
+ schemaId2,
+ schemaId3,
+ },
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId1 = respCreateService.ServiceId
+
+ respPutData, err :=
datasource.Instance().ModifySchema(getContext(), &pb.ModifySchemaRequest{
+ ServiceId: serviceId1,
+ SchemaId: schemaId2,
+ Schema: schemaContent,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respPutData.Response.GetCode())
+
+ respPutData, err =
datasource.Instance().ModifySchema(getContext(), &pb.ModifySchemaRequest{
+ ServiceId: serviceId1,
+ SchemaId: schemaId3,
+ Schema: schemaContent,
+ Summary: summary,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respPutData.Response.GetCode())
+
+ respGetAllSchema, err :=
datasource.Instance().GetAllSchemas(getContext(), &pb.GetAllSchemaRequest{
+ ServiceId: serviceId1,
+ WithSchema: false,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respGetAllSchema.Response.GetCode())
+ schemas := respGetAllSchema.Schemas
+ for _, schema := range schemas {
+ if schema.SchemaId == schemaId1 && schema.SchemaId ==
schemaId2 {
+ assert.Empty(t, schema.Summary)
+ assert.Empty(t, schema.Schema)
+ }
+ if schema.SchemaId == schemaId3 {
+ assert.Equal(t, summary, schema.Summary)
+ assert.Empty(t, schema.Schema)
+ }
+ }
+
+ respGetAllSchema, err =
datasource.Instance().GetAllSchemas(getContext(), &pb.GetAllSchemaRequest{
+ ServiceId: serviceId1,
+ WithSchema: true,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respGetAllSchema.Response.GetCode())
+ schemas = respGetAllSchema.Schemas
+ for _, schema := range schemas {
+ switch schema.SchemaId {
+ case schemaId1:
+ assert.Empty(t, schema.Summary)
+ assert.Empty(t, schema.Schema)
+ case schemaId2:
+ assert.Empty(t, schema.Summary)
+ assert.Equal(t, schemaContent, schema.Schema)
+ case schemaId3:
+ assert.Equal(t, summary, schema.Summary)
+ assert.Equal(t, schemaContent, schema.Schema)
+ }
+ }
+ })
+
+ t.Run("test get when request is invalid", func(t *testing.T) {
+ log.Info("service does not exist")
+ respGetSchema, err :=
datasource.Instance().GetSchema(getContext(), &pb.GetSchemaRequest{
+ ServiceId: "none_exist_service",
+ SchemaId: "com.huawei.test",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
respGetSchema.Response.GetCode())
+
+ respGetAllSchemas, err :=
datasource.Instance().GetAllSchemas(getContext(), &pb.GetAllSchemaRequest{
+ ServiceId: "none_exist_service",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
respGetAllSchemas.Response.GetCode())
+
+ log.Info("schema id doest not exist")
+ respGetSchema, err =
datasource.Instance().GetSchema(getContext(), &pb.GetSchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "none_exist_schema",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrSchemaNotExists,
respGetSchema.Response.GetCode())
+ })
+
+ t.Run("test get when request is valid", func(t *testing.T) {
+ resp, err := datasource.Instance().GetSchema(getContext(),
&pb.GetSchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "com.huawei.test.ms",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ assert.Equal(t, "get schema ms", resp.Schema)
+ assert.Equal(t, "schema0summary1ms", resp.SchemaSummary)
+
+ })
+}
+
+func TestSchema_Delete(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("register service and instance", func(t *testing.T) {
+ respCreateService, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "delete_schema_group_ms",
+ ServiceName: "delete_schema_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId = respCreateService.ServiceId
+
+ resp, err := datasource.Instance().ModifySchema(getContext(),
&pb.ModifySchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "com.huawei.test.ms",
+ Schema: "delete schema ms",
+ Summary: "summary_ms",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+
+ t.Run("test delete when request is invalid", func(t *testing.T) {
+ log.Info("schema id does not exist")
+ resp, err := datasource.Instance().DeleteSchema(getContext(),
&pb.DeleteSchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "none_exist_schema",
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+
+ log.Info("service id does not exist")
+ resp, err = datasource.Instance().DeleteSchema(getContext(),
&pb.DeleteSchemaRequest{
+ ServiceId: "not_exist_service",
+ SchemaId: "com.huawei.test.ms",
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+ })
+
+ t.Run("test delete when request is valid", func(t *testing.T) {
+ resp, err := datasource.Instance().DeleteSchema(getContext(),
&pb.DeleteSchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "com.huawei.test.ms",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ respGet, err := datasource.Instance().GetSchema(getContext(),
&pb.GetSchemaRequest{
+ ServiceId: serviceId,
+ SchemaId: "com.huawei.test.ms",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrSchemaNotExists,
respGet.Response.GetCode())
+
+ respExist, err :=
datasource.Instance().ExistSchema(getContext(), &pb.GetExistenceRequest{
+ Type: "schema",
+ ServiceId: serviceId,
+ SchemaId: "com.huawei.test.ms",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrSchemaNotExists,
respExist.Response.GetCode())
+ })
+}
+
+func TestRule_Add(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 (
+ serviceId1 string
+ serviceId2 string
+ )
+
+ t.Run("register service and instance", func(t *testing.T) {
+ respCreateService, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "create_rule_group_ms",
+ ServiceName: "create_rule_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId1 = respCreateService.ServiceId
+
+ respCreateService, err =
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "create_rule_group_ms",
+ ServiceName: "create_rule_service_ms",
+ Version: "1.0.1",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId2 = respCreateService.ServiceId
+ })
+
+ t.Run("invalid request", func(t *testing.T) {
+ log.Info("service does not exist")
+ respAddRule, err := datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: "not_exist_service_ms",
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test white",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS, respAddRule)
+ })
+
+ t.Run("request is valid", func(t *testing.T) {
+ log.Info("create a new black list")
+ respAddRule, err := datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId1,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test black",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+ ruleId := respAddRule.RuleIds[0]
+ assert.NotEqual(t, "", ruleId)
+
+ log.Info("create the black list again")
+ respAddRule, err = datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId1,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test change black",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+ assert.Equal(t, 0, len(respAddRule.RuleIds))
+
+ log.Info("create a new white list when black list already
exists")
+ respAddRule, err = datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId1,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "WHITE",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test white",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+ })
+
+ t.Run("create rule out of gaugue", func(t *testing.T) {
+ size := quota.DefaultRuleQuota + 1
+ rules := make([]*pb.AddOrUpdateServiceRule, 0, size)
+ for i := 0; i < size; i++ {
+ rules = append(rules, &pb.AddOrUpdateServiceRule{
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: strconv.Itoa(i),
+ Description: "test white",
+ })
+ }
+
+ resp, err := datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId2,
+ Rules: rules[:size-1],
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ resp, err = datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId2,
+ Rules: rules[size-1:],
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrNotEnoughQuota,
resp.Response.GetCode())
+ })
+}
+
+func TestRule_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)
+
+ var (
+ serviceId string
+ ruleId string
+ )
+
+ t.Run("register service and rules", func(t *testing.T) {
+ respCreateService, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "get_rule_group_ms",
+ ServiceName: "get_rule_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId = respCreateService.ServiceId
+
+ respAddRule, err := datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test BLACK",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+ ruleId = respAddRule.RuleIds[0]
+ assert.NotEqual(t, "", ruleId)
+ })
+
+ t.Run("get when request is invalid", func(t *testing.T) {
+ log.Info("service not exists")
+ respGetRule, err := datasource.Instance().GetRule(getContext(),
&pb.GetServiceRulesRequest{
+ ServiceId: "not_exist_service_ms",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
respGetRule.Response.GetCode())
+ })
+
+ t.Run("get when request is valid", func(t *testing.T) {
+ respGetRule, err := datasource.Instance().GetRule(getContext(),
&pb.GetServiceRulesRequest{
+ ServiceId: serviceId,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respGetRule.Response.GetCode())
+ assert.Equal(t, ruleId, respGetRule.Rules[0].RuleId)
+ })
+}
+
+func TestRule_Update(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
+ ruleId string
+ )
+
+ t.Run("create service and rules", func(t *testing.T) {
+ respCreateService, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "update_rule_group_ms",
+ ServiceName: "update_rule_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId = respCreateService.ServiceId
+
+ respAddRule, err := datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test BLACK",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+ ruleId = respAddRule.RuleIds[0]
+ assert.NotEqual(t, "", ruleId)
+ })
+
+ t.Run("update when request is invalid", func(t *testing.T) {
+ rule := &pb.AddOrUpdateServiceRule{
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test BLACK update",
+ }
+ log.Info("service does not exist")
+ resp, err := datasource.Instance().UpdateRule(getContext(),
&pb.UpdateServiceRuleRequest{
+ ServiceId: "not_exist_service_ms",
+ RuleId: ruleId,
+ Rule: rule,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+
+ log.Info("rule not exists")
+ resp, err = datasource.Instance().UpdateRule(getContext(),
&pb.UpdateServiceRuleRequest{
+ ServiceId: serviceId,
+ RuleId: "not_exist_rule_ms",
+ Rule: rule,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+
+ log.Info("change rule type")
+ resp, err = datasource.Instance().UpdateRule(getContext(),
&pb.UpdateServiceRuleRequest{
+ ServiceId: serviceId,
+ RuleId: ruleId,
+ Rule: &pb.AddOrUpdateServiceRule{
+ RuleType: "WHITE",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test white update",
+ },
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+ })
+
+ t.Run("update when request is valid", func(t *testing.T) {
+ resp, err := datasource.Instance().UpdateRule(getContext(),
&pb.UpdateServiceRuleRequest{
+ ServiceId: serviceId,
+ RuleId: ruleId,
+ Rule: &pb.AddOrUpdateServiceRule{
+ RuleType: "BLACK",
+ Attribute: "AppId",
+ Pattern: "Test*",
+ Description: "test white update",
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+}
+
+func TestRule_Delete(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
+ ruleId string
+ )
+ t.Run("register service and rules", func(t *testing.T) {
+ respCreateService, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "delete_rule_group_ms",
+ ServiceName: "delete_rule_service_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId = respCreateService.ServiceId
+
+ respAddRule, err := datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: serviceId,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "BLACK",
+ Attribute: "ServiceName",
+ Pattern: "Test*",
+ Description: "test BLACK",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+ ruleId = respAddRule.RuleIds[0]
+ assert.NotEqual(t, "", ruleId)
+ })
+
+ t.Run("delete when request is invalid", func(t *testing.T) {
+ log.Info("service not exist")
+ resp, err := datasource.Instance().DeleteRule(getContext(),
&pb.DeleteServiceRulesRequest{
+ ServiceId: "not_exist_service_ms",
+ RuleIds: []string{"1000000"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("rule not exist")
+ resp, err = datasource.Instance().DeleteRule(getContext(),
&pb.DeleteServiceRulesRequest{
+ ServiceId: serviceId,
+ RuleIds: []string{"not_exist_rule"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrRuleNotExists, resp.Response.GetCode())
+ })
+
+ t.Run("delete when request is valid", func(t *testing.T) {
+ resp, err := datasource.Instance().DeleteRule(getContext(),
&pb.DeleteServiceRulesRequest{
+ ServiceId: serviceId,
+ RuleIds: []string{ruleId},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ respGetRule, err := datasource.Instance().GetRule(getContext(),
&pb.GetServiceRulesRequest{
+ ServiceId: serviceId,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ assert.Equal(t, 0, len(respGetRule.Rules))
+ })
+}
+
+func TestRule_Permission(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 (
+ consumerVersion string
+ consumerTag string
+ providerBlack string
+ providerWhite string
+ )
+ t.Run("register service and rules", func(t *testing.T) {
+ respCreateService, err :=
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_tag_ms",
+ ServiceName:
"query_instance_version_consumer_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ consumerVersion = respCreateService.ServiceId
+
+ respCreateService, err =
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_tag_ms",
+ ServiceName: "query_instance_tag_service_ms",
+ Version: "1.0.2",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ providerBlack = respCreateService.ServiceId
+
+ respAddRule, err := datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: providerBlack,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "BLACK",
+ Attribute: "Version",
+ Pattern: "1.0.0",
+ },
+ {
+ RuleType: "BLACK",
+ Attribute: "tag_a",
+ Pattern: "b",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+
+ respCreateService, err =
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_tag_ms",
+ ServiceName: "query_instance_tag_service_ms",
+ Version: "1.0.3",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ providerWhite = respCreateService.ServiceId
+
+ respAddRule, err = datasource.Instance().AddRule(getContext(),
&pb.AddServiceRulesRequest{
+ ServiceId: providerWhite,
+ Rules: []*pb.AddOrUpdateServiceRule{
+ {
+ RuleType: "WHITE",
+ Attribute: "Version",
+ Pattern: "1.0.0",
+ },
+ {
+ RuleType: "WHITE",
+ Attribute: "tag_a",
+ Pattern: "b",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddRule.Response.GetCode())
+
+ respCreateService, err =
datasource.Instance().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_tag_ms",
+ ServiceName: "query_instance_tag_consumer_ms",
+ Version: "1.0.4",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ consumerTag = respCreateService.ServiceId
+
+ respAddTag, err := datasource.Instance().AddTags(getContext(),
&pb.AddServiceTagsRequest{
+ ServiceId: consumerTag,
+ Tags: map[string]string{"a": "b"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respAddTag.Response.GetCode())
+ })
+
+ t.Run("when query instances", func(t *testing.T) {
+ log.Info("consumer version in black list")
+ resp, err := datasource.Instance().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: consumerVersion,
+ ProviderServiceId: providerBlack,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("consumer tag in black list")
+ resp, err = datasource.Instance().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: consumerTag,
+ ProviderServiceId: providerBlack,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("find should return 200 even if consumer permission
deny")
+ respFind, err :=
datasource.Instance().FindInstances(getContext(), &pb.FindInstancesRequest{
+ ConsumerServiceId: consumerVersion,
+ AppId: "query_instance_tag_ms",
+ ServiceName: "query_instance_tag_service_ms",
+ VersionRule: "0+",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 0, len(respFind.Instances))
+
+ respFind, err =
datasource.Instance().FindInstances(getContext(), &pb.FindInstancesRequest{
+ ConsumerServiceId: consumerTag,
+ AppId: "query_instance_tag_ms",
+ ServiceName: "query_instance_tag_service_ms",
+ VersionRule: "0+",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 0, len(respFind.Instances))
+
+ log.Info("consumer not in black list")
+ resp, err = datasource.Instance().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: providerWhite,
+ ProviderServiceId: providerBlack,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ log.Info("consumer not in white list")
+ resp, err = datasource.Instance().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: providerBlack,
+ ProviderServiceId: providerWhite,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
resp.Response.GetCode())
+
+ log.Info("consumer version in white list")
+ resp, err = datasource.Instance().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: consumerVersion,
+ ProviderServiceId: providerWhite,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ log.Info("consumer tag in white list")
+ resp, err = datasource.Instance().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: consumerTag,
+ ProviderServiceId: providerWhite,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+}
+
+func TestTags_Add(t *testing.T) {
var (
serviceId1 string
serviceId2 string
)
+
// init
datasource.Install("etcd", func(opts datasource.Options)
(datasource.DataSource, error) {
return NewDataSource(opts), nil
@@ -2790,11 +3604,12 @@ func TestTag_Add(t *testing.T) {
PluginImplName:
datasource.ImplName(archaius.GetString("servicecomb.datasource.name", "etcd")),
})
assert.NoError(t, err)
+
// create service
t.Run("create service", func(t *testing.T) {
svc1 := &pb.MicroService{
- AppId: "create_tag_group",
- ServiceName: "create_tag_service",
+ AppId: "create_tag_group_ms",
+ ServiceName: "create_tag_service_ms",
Version: "1.0.0",
Level: "FRONT",
Status: pb.MS_UP,
@@ -2807,8 +3622,8 @@ func TestTag_Add(t *testing.T) {
serviceId1 = resp.ServiceId
svc2 := &pb.MicroService{
- AppId: "create_tag_group",
- ServiceName: "create_tag_service",
+ AppId: "create_tag_group_ms",
+ ServiceName: "create_tag_service_ms",
Version: "1.0.1",
Level: "FRONT",
Status: pb.MS_UP,
diff --git a/datasource/ms.go b/datasource/ms.go
index 8f07e30..c7bd219 100644
--- a/datasource/ms.go
+++ b/datasource/ms.go
@@ -51,16 +51,17 @@ type MetadataManager interface {
ModifySchemas(ctx context.Context, request *pb.ModifySchemasRequest)
(*pb.ModifySchemasResponse, error)
ModifySchema(ctx context.Context, request *pb.ModifySchemaRequest)
(*pb.ModifySchemaResponse, error)
ExistSchema(ctx context.Context, request *pb.GetExistenceRequest)
(*pb.GetExistenceResponse, error)
- GetSchema()
- DeleteSchema()
+ GetSchema(ctx context.Context, request *pb.GetSchemaRequest)
(*pb.GetSchemaResponse, error)
+ 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()
- AddRule()
- GetRule()
- UpdateRule()
- DeleteRule()
+ AddRule(ctx context.Context, request *pb.AddServiceRulesRequest)
(*pb.AddServiceRulesResponse, error)
+ GetRule(ctx context.Context, request *pb.GetServiceRulesRequest)
(*pb.GetServiceRulesResponse, error)
+ UpdateRule(ctx context.Context, request *pb.UpdateServiceRuleRequest)
(*pb.UpdateServiceRuleResponse, error)
+ DeleteRule(ctx context.Context, request *pb.DeleteServiceRulesRequest)
(*pb.DeleteServiceRulesResponse, error)
}
diff --git a/scripts/ut_test_in_docker.sh b/scripts/ut_test_in_docker.sh
index 3dc640b..d480c60 100755
--- a/scripts/ut_test_in_docker.sh
+++ b/scripts/ut_test_in_docker.sh
@@ -40,6 +40,7 @@ echo "${green}Etcd is running......${reset}"
echo "${green}Preparing the env for UT....${reset}"
./scripts/prepare_env_ut.sh
+[ $? == 0 ] && ut_for_dir datasource
[ $? == 0 ] && ut_for_dir pkg
[ $? == 0 ] && ut_for_dir server
[ $? == 0 ] && ut_for_dir scctl