This is an automated email from the ASF dual-hosted git repository.
littlecui pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/servicecomb-service-center.git
The following commit(s) were added to refs/heads/master by this push:
new 57320fa [enhancement] implement ms module interface (#710)
57320fa is described below
commit 57320faa33ddcb4d9bc42d1f2cd1ca91a2652861
Author: popozy <[email protected]>
AuthorDate: Thu Oct 15 09:25:17 2020 +0800
[enhancement] implement ms module interface (#710)
1. implement service/ms/etcd interface
+ GetInstance
+ GetInstances
+ FindInstances
+ UpdateInstanceStatus
+ UpdateInstanceProperties
+ UnregisterInstance
+ HeartbeatSet
+ BatchFind
2. finish unit test relative to the instance above
---
server/service/ms/datasource.go | 17 +-
server/service/ms/etcd/etcd.go | 580 +++++++++++++++++-
server/service/ms/etcd/etcd_test.go | 1148 +++++++++++++++++++++++++++++++++++
server/service/ms/etcd/util.go | 36 ++
4 files changed, 1768 insertions(+), 13 deletions(-)
diff --git a/server/service/ms/datasource.go b/server/service/ms/datasource.go
index 1c93b82..abf890a 100644
--- a/server/service/ms/datasource.go
+++ b/server/service/ms/datasource.go
@@ -31,10 +31,19 @@ type DataSource interface {
GetDeleteServiceFunc(ctx context.Context, serviceID string, force bool,
serviceRespChan chan<- *pb.DelServicesRspInfo)
func(context.Context)
- RegisterInstance(ctx context.Context, in *pb.RegisterInstanceRequest)
(*pb.RegisterInstanceResponse, error)
- SearchInstance()
- UpdateInstance()
- UnRegisterInstance()
+ RegisterInstance(ctx context.Context, request
*pb.RegisterInstanceRequest) (*pb.RegisterInstanceResponse, error)
+ GetInstance(ctx context.Context, request *pb.GetOneInstanceRequest)
(*pb.GetOneInstanceResponse, error)
+ GetInstances(ctx context.Context, request *pb.GetInstancesRequest)
(*pb.GetInstancesResponse, error)
+ FindInstances(ctx context.Context, request *pb.FindInstancesRequest)
(*pb.FindInstancesResponse, error)
+ UpdateInstanceStatus(ctx context.Context, request
*pb.UpdateInstanceStatusRequest) (
+ *pb.UpdateInstanceStatusResponse, error)
+ UpdateInstanceProperties(ctx context.Context, request
*pb.UpdateInstancePropsRequest) (
+ *pb.UpdateInstancePropsResponse, error)
+ UnregisterInstance(ctx context.Context, request
*pb.UnregisterInstanceRequest) (*pb.UnregisterInstanceResponse,
+ error)
+ Heartbeat(ctx context.Context, request *pb.HeartbeatRequest)
(*pb.HeartbeatResponse, error)
+ HeartbeatSet(ctx context.Context, request *pb.HeartbeatSetRequest)
(*pb.HeartbeatSetResponse, error)
+ BatchFind(ctx context.Context, request *pb.BatchFindInstancesRequest)
(*pb.BatchFindInstancesResponse, error)
ModifySchemas(ctx context.Context, request *pb.ModifySchemasRequest)
(*pb.ModifySchemasResponse, error)
ModifySchema(ctx context.Context, request *pb.ModifySchemaRequest)
(*pb.ModifySchemaResponse, error)
diff --git a/server/service/ms/etcd/etcd.go b/server/service/ms/etcd/etcd.go
index aba1669..31f7b3a 100644
--- a/server/service/ms/etcd/etcd.go
+++ b/server/service/ms/etcd/etcd.go
@@ -20,6 +20,7 @@ import (
"encoding/json"
"errors"
"fmt"
+ "github.com/apache/servicecomb-service-center/pkg/gopool"
"github.com/apache/servicecomb-service-center/pkg/log"
pb "github.com/apache/servicecomb-service-center/pkg/registry"
"github.com/apache/servicecomb-service-center/pkg/util"
@@ -30,6 +31,7 @@ import (
"github.com/apache/servicecomb-service-center/server/plugin/quota"
"github.com/apache/servicecomb-service-center/server/plugin/registry"
scerr "github.com/apache/servicecomb-service-center/server/scerror"
+ "github.com/apache/servicecomb-service-center/server/service/cache"
"github.com/apache/servicecomb-service-center/server/service/ms"
serviceUtil
"github.com/apache/servicecomb-service-center/server/service/util"
"strconv"
@@ -381,16 +383,573 @@ func (ds *DataSource) RegisterInstance(ctx
context.Context, request *pb.Register
}, nil
}
-func (ds *DataSource) SearchInstance() {
- panic("implement me")
+func (ds *DataSource) GetInstance(ctx context.Context, in
*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 err != nil {
+ log.Errorf(err, "get consumer failed, consumer[%s] find
provider instance[%s/%s]",
+ in.ConsumerServiceId, in.ProviderServiceId,
in.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)
+ return &pb.GetOneInstanceResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
+ fmt.Sprintf("Consumer[%s] does not
exist.", in.ConsumerServiceId)),
+ }, nil
+ }
+ }
+
+ provider, err := serviceUtil.GetService(ctx, domainProject,
in.ProviderServiceId)
+ if err != nil {
+ log.Errorf(err, "get provider failed, consumer[%s] find
provider instance[%s/%s]",
+ in.ConsumerServiceId, in.ProviderServiceId,
in.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)
+ return &pb.GetOneInstanceResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
+ fmt.Sprintf("Provider[%s] does not exist.",
in.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,
+ provider.ServiceId, provider.Environment,
provider.AppId, provider.ServiceName, provider.Version,
+ in.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)
+ if err != nil {
+ log.Errorf(err, "FindInstances.GetWithProviderID failed, %s
failed", findFlag())
+ return &pb.GetOneInstanceResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if item == nil || len(item.Instances) == 0 {
+ mes := fmt.Errorf("%s failed, provider instance does not
exist", findFlag())
+ log.Errorf(mes, "FindInstances.GetWithProviderID failed")
+ return &pb.GetOneInstanceResponse{
+ Response:
proto.CreateResponse(scerr.ErrInstanceNotExists, mes.Error()),
+ }, nil
+ }
+
+ instance := item.Instances[0]
+ if rev == item.Rev {
+ instance = nil // for gRPC
+ }
+ _ = util.SetContext(ctx, util.CtxResponseRevision, item.Rev)
+
+ return &pb.GetOneInstanceResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Get
instance successfully."),
+ Instance: instance,
+ }, nil
}
-func (ds *DataSource) UpdateInstance() {
- panic("implement me")
+func (ds *DataSource) FindInstances(ctx context.Context, in
*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,
+ }
+
+ rev, ok := ctx.Value(util.CtxRequestRevision).(string)
+ if !ok {
+ err := errors.New("rev in context is not type string")
+ log.Error("", err)
+ return &pb.FindInstancesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ if apt.IsShared(provider) {
+ return ds.findSharedServiceInstance(ctx, in, provider, rev)
+ }
+ return ds.findInstance(ctx, in, provider, rev)
}
-func (ds *DataSource) UnRegisterInstance() {
- panic("implement me")
+func (ds *DataSource) findInstance(ctx context.Context, in
*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)
+ 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)
+ 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)
+ return &pb.FindInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
+ fmt.Sprintf("Consumer[%s] does not
exist.", in.ConsumerServiceId)),
+ }, nil
+ }
+ provider.Environment = service.Environment
+ }
+
+ // provider is not a shared micro-service,
+ // only allow shared micro-service instances found in 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,
+ provider.Environment, provider.AppId, provider.ServiceName,
provider.Version)
+
+ // cache
+ var item *cache.VersionRuleCacheItem
+ item, err = cache.FindInstances.Get(ctx, service, provider, in.Tags,
rev)
+ if err != nil {
+ log.Errorf(err, "FindInstancesCache.Get failed, %s failed",
findFlag)
+ return &pb.FindInstancesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if item == nil {
+ mes := fmt.Errorf("%s failed, provider does not exist",
findFlag)
+ log.Errorf(mes, "FindInstancesCache.Get failed")
+ return &pb.FindInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, mes.Error()),
+ }, nil
+ }
+
+ // add dependency queue
+ if len(in.ConsumerServiceId) > 0 &&
+ len(item.ServiceIds) > 0 &&
+ !cache.DependencyRule.ExistVersionRule(ctx,
in.ConsumerServiceId, provider) {
+ provider, err = ds.reshapeProviderKey(ctx, provider,
item.ServiceIds[0])
+ if err != nil {
+ return nil, err
+ }
+ if provider != nil {
+ err = serviceUtil.AddServiceVersionRule(ctx,
domainProject, service, provider)
+ } else {
+ mes := fmt.Errorf("%s failed, provider does not exist",
findFlag)
+ log.Errorf(mes, "AddServiceVersionRule failed")
+ return &pb.FindInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, mes.Error()),
+ }, nil
+ }
+ if err != nil {
+ log.Errorf(err, "AddServiceVersionRule failed, %s
failed", findFlag)
+ return &pb.FindInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrInternal, err.Error()),
+ }, err
+ }
+ }
+
+ return ds.genFindResult(ctx, rev, item)
+}
+
+func (ds *DataSource) findSharedServiceInstance(ctx context.Context, in
*pb.FindInstancesRequest,
+ provider *pb.MicroServiceKey, rev string) (*pb.FindInstancesResponse,
error) {
+ var err error
+ service := &pb.MicroService{Environment: in.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)
+ if err != nil {
+ log.Errorf(err, "FindInstancesCache.Get failed, %s failed",
findFlag)
+ return &pb.FindInstancesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if item == nil {
+ mes := fmt.Errorf("%s failed, provider does not exist",
findFlag)
+ log.Errorf(mes, "FindInstancesCache.Get failed")
+ return &pb.FindInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, mes.Error()),
+ }, nil
+ }
+
+ return ds.genFindResult(ctx, rev, item)
+}
+
+func (ds *DataSource) genFindResult(ctx context.Context, oldRev string, item
*cache.VersionRuleCacheItem) (
+ *pb.FindInstancesResponse, error) {
+ instances := item.Instances
+ if oldRev == item.Rev {
+ instances = nil // for gRPC
+ }
+ // TODO support gRPC output context
+ _ = util.SetContext(ctx, util.CtxResponseRevision, item.Rev)
+ return &pb.FindInstancesResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Query
service instances successfully."),
+ Instances: instances,
+ }, nil
+}
+
+func (ds *DataSource) reshapeProviderKey(ctx context.Context, provider
*pb.MicroServiceKey, providerID string) (
+ *pb.MicroServiceKey, error) {
+ //维护version的规则,service name 可能是别名,所以重新获取
+ providerService, err := serviceUtil.GetService(ctx, provider.Tenant,
providerID)
+ if providerService == nil {
+ return nil, err
+ }
+
+ versionRule := provider.Version
+ provider = proto.MicroServiceToKey(provider.Tenant, providerService)
+ provider.Version = versionRule
+ return provider, nil
+}
+
+func (ds *DataSource) UpdateInstanceStatus(ctx context.Context, in
*pb.UpdateInstanceStatusRequest) (*pb.
+ UpdateInstanceStatusResponse, error) {
+ domainProject := util.ParseDomainProject(ctx)
+ updateStatusFlag := util.StringJoin([]string{in.ServiceId,
in.InstanceId, in.Status}, "/")
+
+ instance, err := serviceUtil.GetInstance(ctx, domainProject,
in.ServiceId, in.InstanceId)
+ if err != nil {
+ log.Errorf(err, "update instance[%s] status failed",
updateStatusFlag)
+ return &pb.UpdateInstanceStatusResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if instance == nil {
+ log.Errorf(nil, "update instance[%s] status failed, instance
does not exist", updateStatusFlag)
+ return &pb.UpdateInstanceStatusResponse{
+ Response:
proto.CreateResponse(scerr.ErrInstanceNotExists, "Service instance does not
exist."),
+ }, nil
+ }
+
+ copyInstanceRef := *instance
+ copyInstanceRef.Status = in.Status
+
+ if err := serviceUtil.UpdateInstance(ctx, domainProject,
©InstanceRef); err != nil {
+ log.Errorf(err, "update instance[%s] status failed",
updateStatusFlag)
+ resp := &pb.UpdateInstanceStatusResponse{
+ Response: proto.CreateResponseWithSCErr(err),
+ }
+ if err.InternalError() {
+ return resp, err
+ }
+ return resp, nil
+ }
+
+ log.Infof("update instance[%s] status successfully", updateStatusFlag)
+ return &pb.UpdateInstanceStatusResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Update
service instance status successfully."),
+ }, nil
+}
+
+func (ds *DataSource) UpdateInstanceProperties(ctx context.Context, in
*pb.UpdateInstancePropsRequest) (
+ *pb.UpdateInstancePropsResponse, error) {
+ domainProject := util.ParseDomainProject(ctx)
+ instanceFlag := util.StringJoin([]string{in.ServiceId, in.InstanceId},
"/")
+
+ instance, err := serviceUtil.GetInstance(ctx, domainProject,
in.ServiceId, in.InstanceId)
+ if err != nil {
+ log.Errorf(err, "update instance[%s] properties failed",
instanceFlag)
+ return &pb.UpdateInstancePropsResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if instance == nil {
+ log.Errorf(nil, "update instance[%s] properties failed,
instance does not exist", instanceFlag)
+ return &pb.UpdateInstancePropsResponse{
+ Response:
proto.CreateResponse(scerr.ErrInstanceNotExists, "Service instance does not
exist."),
+ }, nil
+ }
+
+ copyInstanceRef := *instance
+ copyInstanceRef.Properties = in.Properties
+
+ if err := serviceUtil.UpdateInstance(ctx, domainProject,
©InstanceRef); err != nil {
+ log.Errorf(err, "update instance[%s] properties failed",
instanceFlag)
+ resp := &pb.UpdateInstancePropsResponse{
+ Response: proto.CreateResponseWithSCErr(err),
+ }
+ if err.InternalError() {
+ return resp, err
+ }
+ return resp, nil
+ }
+
+ log.Infof("update instance[%s] properties successfully", instanceFlag)
+ return &pb.UpdateInstancePropsResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Update
service instance properties successfully."),
+ }, nil
+}
+
+func (ds *DataSource) HeartbeatSet(ctx context.Context, in
*pb.HeartbeatSetRequest) (*pb.HeartbeatSetResponse, error) {
+ domainProject := util.ParseDomainProject(ctx)
+
+ heartBeatCount := len(in.Instances)
+ existFlag := make(map[string]bool, heartBeatCount)
+ instancesHbRst := make(chan *pb.InstanceHbRst, heartBeatCount)
+ noMultiCounter := 0
+ for _, heartbeatElement := range in.Instances {
+ if _, ok :=
existFlag[heartbeatElement.ServiceId+heartbeatElement.InstanceId]; ok {
+ log.Warnf("instance[%s/%s] is duplicate in heartbeat
set",
+ heartbeatElement.ServiceId,
heartbeatElement.InstanceId)
+ continue
+ } else {
+
existFlag[heartbeatElement.ServiceId+heartbeatElement.InstanceId] = true
+ noMultiCounter++
+ }
+ gopool.Go(getHeartbeatFunc(ctx, domainProject, instancesHbRst,
heartbeatElement))
+ }
+ count := 0
+ successFlag := false
+ failFlag := false
+ instanceHbRstArr := make([]*pb.InstanceHbRst, 0, heartBeatCount)
+ for heartbeat := range instancesHbRst {
+ count++
+ if len(heartbeat.ErrMessage) != 0 {
+ failFlag = true
+ } else {
+ successFlag = true
+ }
+ instanceHbRstArr = append(instanceHbRstArr, heartbeat)
+ if count == noMultiCounter {
+ close(instancesHbRst)
+ }
+ }
+ if !failFlag && successFlag {
+ log.Infof("batch update heartbeats[%s] successfully", count)
+ return &pb.HeartbeatSetResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
"Heartbeat set successfully."),
+ Instances: instanceHbRstArr,
+ }, nil
+ }
+ log.Errorf(nil, "batch update heartbeats failed, %v", in.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,
+ 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 err != nil {
+ log.Errorf(err, "get consumer failed, consumer[%s] find
provider instances",
+ in.ConsumerServiceId, in.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)
+ return &pb.GetInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
+ fmt.Sprintf("Consumer[%s] does not
exist.", in.ConsumerServiceId)),
+ }, nil
+ }
+ }
+
+ provider, err := serviceUtil.GetService(ctx, domainProject,
in.ProviderServiceId)
+ if err != nil {
+ log.Errorf(err, "get provider failed, consumer[%s] find
provider instances",
+ in.ConsumerServiceId, in.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)
+ return &pb.GetInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists,
+ fmt.Sprintf("Provider[%s] does not exist.",
in.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,
+ provider.ServiceId, provider.Environment,
provider.AppId, provider.ServiceName, provider.Version)
+ }
+
+ 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,
+ }, in.Tags, rev)
+ if err != nil {
+ log.Errorf(err, "FindInstances.GetWithProviderID failed, %s
failed", findFlag())
+ return &pb.GetInstancesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+ if item == nil || len(item.ServiceIds) == 0 {
+ mes := fmt.Errorf("%s failed, provider instance does not
exist", findFlag())
+ log.Errorf(mes, "FindInstances.GetWithProviderID failed")
+ return &pb.GetInstancesResponse{
+ Response:
proto.CreateResponse(scerr.ErrServiceNotExists, mes.Error()),
+ }, nil
+ }
+
+ instances := item.Instances
+ if rev == item.Rev {
+ instances = nil // for gRPC
+ }
+ _ = util.SetContext(ctx, util.CtxResponseRevision, item.Rev)
+
+ return &pb.GetInstancesResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Query
service instances successfully."),
+ Instances: instances,
+ }, nil
+}
+
+func (ds *DataSource) BatchFind(ctx context.Context, request
*pb.BatchFindInstancesRequest) (
+ *pb.BatchFindInstancesResponse, error) {
+ response := &pb.BatchFindInstancesResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS, "Batch
query service instances successfully."),
+ }
+
+ var err error
+ // find services
+ response.Services, err = ds.batchFindServices(ctx, request)
+ if err != nil {
+ return &pb.BatchFindInstancesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ // find instance
+ response.Instances, err = ds.batchFindInstances(ctx, request)
+ if err != nil {
+ return &pb.BatchFindInstancesResponse{
+ Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
+ }, err
+ }
+
+ return response, nil
+}
+
+func (ds *DataSource) batchFindServices(ctx context.Context, in
*pb.BatchFindInstancesRequest) (
+ *pb.BatchFindResult, error) {
+ if len(in.Services) == 0 {
+ return nil, nil
+ }
+ cloneCtx := util.CloneContext(ctx)
+
+ services := &pb.BatchFindResult{}
+ failedResult := make(map[int32]*pb.FindFailedResult)
+ for index, key := range in.Services {
+ findCtx := util.SetContext(cloneCtx, util.CtxRequestRevision,
key.Rev)
+ resp, err := ds.FindInstances(findCtx, &pb.FindInstancesRequest{
+ ConsumerServiceId: in.ConsumerServiceId,
+ AppId: key.Service.AppId,
+ ServiceName: key.Service.ServiceName,
+ VersionRule: key.Service.Version,
+ Environment: key.Service.Environment,
+ })
+ if err != nil {
+ return nil, err
+ }
+ failed, ok := failedResult[resp.Response.GetCode()]
+ serviceUtil.AppendFindResponse(findCtx, int64(index),
resp.Response, resp.Instances,
+ &services.Updated, &services.NotModified, &failed)
+ if !ok && failed != nil {
+ failedResult[resp.Response.GetCode()] = failed
+ }
+ }
+ for _, result := range failedResult {
+ services.Failed = append(services.Failed, result)
+ }
+ return services, nil
+}
+
+func (ds *DataSource) batchFindInstances(ctx context.Context, in
*pb.BatchFindInstancesRequest) (*pb.BatchFindResult, error) {
+ if len(in.Instances) == 0 {
+ return nil, nil
+ }
+ cloneCtx := util.CloneContext(ctx)
+ // can not find the shared provider instances
+ cloneCtx = util.SetTargetDomainProject(cloneCtx, util.ParseDomain(ctx),
util.ParseProject(ctx))
+
+ instances := &pb.BatchFindResult{}
+ failedResult := make(map[int32]*pb.FindFailedResult)
+ for index, key := range in.Instances {
+ getCtx := util.SetContext(cloneCtx, util.CtxRequestRevision,
key.Rev)
+ resp, err := ds.GetInstance(getCtx, &pb.GetOneInstanceRequest{
+ ConsumerServiceId: in.ConsumerServiceId,
+ ProviderServiceId: key.Instance.ServiceId,
+ ProviderInstanceId: key.Instance.InstanceId,
+ })
+ if err != nil {
+ return nil, err
+ }
+ failed, ok := failedResult[resp.Response.GetCode()]
+ serviceUtil.AppendFindResponse(getCtx, int64(index),
resp.Response, []*pb.MicroServiceInstance{resp.Instance},
+ &instances.Updated, &instances.NotModified, &failed)
+ if !ok && failed != nil {
+ failedResult[resp.Response.GetCode()] = failed
+ }
+ }
+ for _, result := range failedResult {
+ instances.Failed = append(instances.Failed, result)
+ }
+ return instances, nil
+}
+
+func (ds *DataSource) UnregisterInstance(ctx context.Context, request
*pb.UnregisterInstanceRequest) (
+ *pb.UnregisterInstanceResponse, error) {
+ remoteIP := util.GetIPFromContext(ctx)
+ domainProject := util.ParseDomainProject(ctx)
+ serviceID := request.ServiceId
+ instanceID := request.InstanceId
+
+ instanceFlag := util.StringJoin([]string{serviceID, instanceID}, "/")
+
+ err := revokeInstance(ctx, domainProject, serviceID, instanceID)
+ if err != nil {
+ log.Errorf(err, "unregister instance failed, instance[%s],
operator %s: revoke instance failed",
+ instanceFlag, remoteIP)
+ resp := &pb.UnregisterInstanceResponse{
+ Response: proto.CreateResponseWithSCErr(err),
+ }
+ if err.InternalError() {
+ return resp, err
+ }
+ return resp, nil
+ }
+
+ log.Infof("unregister instance[%s], operator %s", instanceFlag,
remoteIP)
+ return &pb.UnregisterInstanceResponse{
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
"Unregister service instance successfully."),
+ }, nil
}
func (ds *DataSource) Heartbeat(ctx context.Context, request
*pb.HeartbeatRequest) (*pb.HeartbeatResponse, error) {
@@ -415,10 +974,12 @@ func (ds *DataSource) Heartbeat(ctx context.Context,
request *pb.HeartbeatReques
log.Errorf(errors.New("connect backend timed out"),
"heartbeat successful, but renew instance[%s] failed.
operator %s", instanceFlag, remoteIP)
} else {
- log.Infof("heartbeat successful, renew instance[%s] ttl to %d.
operator %s", instanceFlag, ttl, remoteIP)
+ log.Infof("heartbeat successful, renew instance[%s] ttl to %d.
operator %s",
+ instanceFlag, ttl, remoteIP)
}
return &pb.HeartbeatResponse{
- Response: proto.CreateResponse(proto.Response_SUCCESS, "Update
service instance heartbeat successfully."),
+ Response: proto.CreateResponse(proto.Response_SUCCESS,
+ "Update service instance heartbeat successfully."),
}, nil
}
@@ -430,7 +991,8 @@ func (ds *DataSource) ModifySchemas(ctx context.Context,
request *pb.ModifySchem
serviceInfo, err := serviceUtil.GetService(ctx, domainProject,
serviceID)
if err != nil {
- log.Errorf(err, "modify service[%s] schemas failed, get service
failed, operator: %s", serviceID, remoteIP)
+ log.Errorf(err, "modify service[%s] schemas failed, get service
failed, operator: %s",
+ serviceID, remoteIP)
return &pb.ModifySchemasResponse{
Response: proto.CreateResponse(scerr.ErrInternal,
err.Error()),
}, err
diff --git a/server/service/ms/etcd/etcd_test.go
b/server/service/ms/etcd/etcd_test.go
index 8811f1f..20b54ff 100644
--- a/server/service/ms/etcd/etcd_test.go
+++ b/server/service/ms/etcd/etcd_test.go
@@ -18,6 +18,8 @@ package etcd_test
import (
"github.com/apache/servicecomb-service-center/pkg/log"
pb "github.com/apache/servicecomb-service-center/pkg/registry"
+ "github.com/apache/servicecomb-service-center/pkg/util"
+ "github.com/apache/servicecomb-service-center/server/core"
"github.com/apache/servicecomb-service-center/server/core/proto"
"github.com/apache/servicecomb-service-center/server/plugin/quota"
scerr "github.com/apache/servicecomb-service-center/server/scerror"
@@ -26,6 +28,7 @@ import (
"github.com/apache/servicecomb-service-center/server/service/ms/etcd"
"github.com/go-chassis/go-archaius"
"github.com/stretchr/testify/assert"
+ "os"
"strconv"
"strings"
"testing"
@@ -817,6 +820,1151 @@ func TestInstance_Create(t *testing.T) {
})
}
+func TestInstance_HeartBeat(t *testing.T) {
+ ms.Install("etcd", func(opts ms.Options) (ms.DataSource, error) {
+ return etcd.NewDataSource(opts), nil
+ })
+ err := ms.Init(ms.Options{
+ Endpoint: "",
+ PluginImplName:
ms.ImplName(archaius.GetString("servicecomb.ms.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ var (
+ serviceId string
+ instanceId1 string
+ instanceId2 string
+ )
+
+ t.Run("register service and instance, should pass", func(t *testing.T) {
+ log.Info("register service")
+ respCreateService, err :=
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ ServiceName: "heartbeat_service_ms",
+ AppId: "heartbeat_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
+
+ respCreateInstance, err :=
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "heartbeat:127.0.0.1:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId1 = respCreateInstance.InstanceId
+
+ respCreateInstance, err =
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "heartbeat:127.0.0.2:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId2 = respCreateInstance.InstanceId
+ })
+
+ t.Run("update a lease", func(t *testing.T) {
+ log.Info("valid instance")
+ resp, err := ms.MicroService().Heartbeat(getContext(),
&pb.HeartbeatRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId1,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ log.Info("serviceId does not exist")
+ resp, err = ms.MicroService().Heartbeat(getContext(),
&pb.HeartbeatRequest{
+ ServiceId: "100000000000",
+ InstanceId: instanceId1,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+
+ log.Info("instance does not exist")
+ resp, err = ms.MicroService().Heartbeat(getContext(),
&pb.HeartbeatRequest{
+ ServiceId: serviceId,
+ InstanceId: "not-exist-ins",
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+ })
+
+ t.Run("batch update lease", func(t *testing.T) {
+ log.Info("request contains at least 1 instances")
+ resp, err := ms.MicroService().HeartbeatSet(getContext(),
&pb.HeartbeatSetRequest{
+ Instances: []*pb.HeartbeatSetElement{
+ {
+ ServiceId: serviceId,
+ InstanceId: instanceId1,
+ },
+ {
+ ServiceId: serviceId,
+ InstanceId: instanceId2,
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+}
+
+func TestInstance_Update(t *testing.T) {
+ ms.Install("etcd", func(opts ms.Options) (ms.DataSource, error) {
+ return etcd.NewDataSource(opts), nil
+ })
+ err := ms.Init(ms.Options{
+ Endpoint: "",
+ PluginImplName:
ms.ImplName(archaius.GetString("servicecomb.ms.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ var (
+ serviceId string
+ instanceId string
+ )
+
+ t.Run("register service and instance, should pass", func(t *testing.T) {
+ log.Info("register service")
+ respCreateService, err :=
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ ServiceName: "update_instance_service_ms",
+ AppId: "update_instance_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
+
+ log.Info("create instance")
+ respCreateInstance, err :=
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId,
+ Endpoints: []string{
+ "updateInstance:127.0.0.1:8080",
+ },
+ HostName: "UT-HOST-MS",
+ Status: pb.MSI_UP,
+ Properties: map[string]string{"nodeIP": "test"},
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId = respCreateInstance.InstanceId
+ })
+
+ t.Run("update instance status", func(t *testing.T) {
+ log.Info("update instance status to DOWN")
+ respUpdateStatus, err :=
ms.MicroService().UpdateInstanceStatus(getContext(),
&pb.UpdateInstanceStatusRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ Status: pb.MSI_DOWN,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateStatus.Response.GetCode())
+
+ log.Info("update instance status to OUTOFSERVICE")
+ respUpdateStatus, err =
ms.MicroService().UpdateInstanceStatus(getContext(),
&pb.UpdateInstanceStatusRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ Status: pb.MSI_OUTOFSERVICE,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateStatus.Response.GetCode())
+
+ log.Info("update instance status to STARTING")
+ respUpdateStatus, err =
ms.MicroService().UpdateInstanceStatus(getContext(),
&pb.UpdateInstanceStatusRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ Status: pb.MSI_STARTING,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateStatus.Response.GetCode())
+
+ log.Info("update instance status to TESTING")
+ respUpdateStatus, err =
ms.MicroService().UpdateInstanceStatus(getContext(),
&pb.UpdateInstanceStatusRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ Status: pb.MSI_TESTING,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateStatus.Response.GetCode())
+
+ log.Info("update instance status to UP")
+ respUpdateStatus, err =
ms.MicroService().UpdateInstanceStatus(getContext(),
&pb.UpdateInstanceStatusRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ Status: pb.MSI_UP,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateStatus.Response.GetCode())
+
+ log.Info("update instance status with a not exist instance")
+ respUpdateStatus, err =
ms.MicroService().UpdateInstanceStatus(getContext(),
&pb.UpdateInstanceStatusRequest{
+ ServiceId: serviceId,
+ InstanceId: "notexistins",
+ Status: pb.MSI_STARTING,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respUpdateStatus.Response.GetCode())
+ })
+
+ t.Run("update instance properties", func(t *testing.T) {
+ log.Info("update one properties")
+ respUpdateProperties, err :=
ms.MicroService().UpdateInstanceProperties(getContext(),
+ &pb.UpdateInstancePropsRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ Properties: map[string]string{
+ "test": "test",
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateProperties.Response.GetCode())
+
+ log.Info("all max properties updated")
+ size := 1000
+ properties := make(map[string]string, size)
+ for i := 0; i < size; i++ {
+ s := strconv.Itoa(i) + strings.Repeat("x", 253)
+ properties[s] = s
+ }
+ respUpdateProperties, err =
ms.MicroService().UpdateInstanceProperties(getContext(),
+ &pb.UpdateInstancePropsRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ Properties: properties,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateProperties.Response.GetCode())
+
+ log.Info("update instance that does not exist")
+ respUpdateProperties, err =
ms.MicroService().UpdateInstanceProperties(getContext(),
+ &pb.UpdateInstancePropsRequest{
+ ServiceId: serviceId,
+ InstanceId: "not_exist_ins",
+ Properties: map[string]string{
+ "test": "test",
+ },
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respUpdateProperties.Response.GetCode())
+
+ log.Info("remove properties")
+ respUpdateProperties, err =
ms.MicroService().UpdateInstanceProperties(getContext(),
+ &pb.UpdateInstancePropsRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respUpdateProperties.Response.GetCode())
+
+ log.Info("update service that does not exist")
+ respUpdateProperties, err =
ms.MicroService().UpdateInstanceProperties(getContext(),
+ &pb.UpdateInstancePropsRequest{
+ ServiceId: "not_exist_service",
+ InstanceId: instanceId,
+ Properties: map[string]string{
+ "test": "test",
+ },
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respUpdateProperties.Response.GetCode())
+ })
+}
+
+func TestInstance_Query(t *testing.T) {
+ ms.Install("etcd", func(opts ms.Options) (ms.DataSource, error) {
+ return etcd.NewDataSource(opts), nil
+ })
+ err := ms.Init(ms.Options{
+ Endpoint: "",
+ PluginImplName:
ms.ImplName(archaius.GetString("servicecomb.ms.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ var (
+ serviceId1 string
+ serviceId2 string
+ serviceId3 string
+ serviceId4 string
+ serviceId5 string
+ serviceId6 string
+ serviceId7 string
+ serviceId8 string
+ serviceId9 string
+ instanceId1 string
+ instanceId2 string
+ instanceId4 string
+ instanceId5 string
+ instanceId8 string
+ instanceId9 string
+ )
+
+ t.Run("register services and instances for testInstance_query", func(t
*testing.T) {
+ respCreateService, err :=
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_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 =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_service_ms",
+ Version: "1.0.5",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId2 = respCreateService.ServiceId
+
+ respCreateService, err =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_diff_app_ms",
+ ServiceName: "query_instance_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())
+ serviceId3 = respCreateService.ServiceId
+
+ respCreateService, err =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ Environment: pb.ENV_PROD,
+ AppId: "query_instance_ms",
+ ServiceName:
"query_instance_diff_env_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())
+ serviceId4 = respCreateService.ServiceId
+
+ respCreateService, err =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ Environment: pb.ENV_PROD,
+ AppId: "default",
+ ServiceName:
"query_instance_shared_provider_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ Properties: map[string]string{
+ proto.PROP_ALLOW_CROSS_APP: "true",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId5 = respCreateService.ServiceId
+
+ respCreateService, err = ms.MicroService().RegisterService(
+ util.SetDomainProject(util.CloneContext(getContext()),
"user", "user"),
+ &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "default",
+ ServiceName:
"query_instance_diff_domain_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())
+ serviceId6 = respCreateService.ServiceId
+
+ respCreateService, err =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "default",
+ ServiceName:
"query_instance_shared_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())
+ serviceId7 = respCreateService.ServiceId
+
+ respCreateService, err =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_with_rev_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId8 = respCreateService.ServiceId
+
+ respCreateService, err =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "query_instance_ms",
+ ServiceName: "batch_query_instance_with_rev_ms",
+ Version: "1.0.0",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId9 = respCreateService.ServiceId
+
+ respCreateInstance, err :=
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId1,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "find:127.0.0.1:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId1 = respCreateInstance.InstanceId
+
+ respCreateInstance, err =
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId2,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "find:127.0.0.2:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId2 = respCreateInstance.InstanceId
+
+ respCreateInstance, err =
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId4,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "find:127.0.0.4:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId4 = respCreateInstance.InstanceId
+
+ respCreateInstance, err =
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId5,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "find:127.0.0.5:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId5 = respCreateInstance.InstanceId
+
+ respCreateInstance, err =
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId8,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "find:127.0.0.8:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId8 = respCreateInstance.InstanceId
+
+ respCreateInstance, err =
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId9,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "find:127.0.0.9:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId9 = respCreateInstance.InstanceId
+ })
+
+ t.Run("query instance", func(t *testing.T) {
+ log.Info("find with version rule")
+ respFind, err := ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId1,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_service_ms",
+ VersionRule: "latest",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, instanceId2, respFind.Instances[0].InstanceId)
+
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId1,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_service_ms",
+ VersionRule: "1.0.0+",
+ Tags: []string{},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, instanceId2, respFind.Instances[0].InstanceId)
+
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId1,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_service_ms",
+ VersionRule: "1.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, instanceId1, respFind.Instances[0].InstanceId)
+
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId1,
+ AppId: "query_instance",
+ ServiceName: "query_instance_service",
+ VersionRule: "0.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrServiceNotExists,
respFind.Response.GetCode())
+
+ log.Info("find with env")
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId4,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_diff_env_service_ms",
+ VersionRule: "1.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Instances))
+ assert.Equal(t, instanceId4, respFind.Instances[0].InstanceId)
+
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ Environment: pb.ENV_PROD,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_diff_env_service_ms",
+ VersionRule: "1.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Instances))
+ assert.Equal(t, instanceId4, respFind.Instances[0].InstanceId)
+
+ log.Info("find with rev")
+ ctx := util.SetContext(getContext(), util.CtxNocache, "")
+ respFind, err = ms.MicroService().FindInstances(ctx,
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId8,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_with_rev_ms",
+ VersionRule: "1.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ rev, _ := ctx.Value(util.CtxResponseRevision).(string)
+ assert.Equal(t, instanceId8, respFind.Instances[0].InstanceId)
+ assert.NotEqual(t, 0, len(rev))
+
+ util.SetContext(ctx, util.CtxRequestRevision, "x")
+ respFind, err = ms.MicroService().FindInstances(ctx,
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId8,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_with_rev_ms",
+ VersionRule: "1.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, instanceId8, respFind.Instances[0].InstanceId)
+ assert.Equal(t, ctx.Value(util.CtxResponseRevision), rev)
+
+ log.Info("find should return 200 if consumer is diff apps")
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId3,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_service_ms",
+ VersionRule: "1.0.5",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 0, len(respFind.Instances))
+
+ log.Info("provider tag does not exist")
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId1,
+ AppId: "query_instance_ms",
+ ServiceName: "query_instance_service_ms",
+ VersionRule: "latest",
+ Tags: []string{"not_exist_tag"},
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 0, len(respFind.Instances))
+
+ log.Info("shared service discovery")
+ _ = os.Setenv("CSE_SHARED_SERVICES",
"query_instance_shared_provider_ms")
+ core.SetSharedMode()
+ core.Service.Environment = pb.ENV_PROD
+ respFind, err = ms.MicroService().FindInstances(
+ util.SetTargetDomainProject(
+
util.SetDomainProject(util.CloneContext(getContext()), "user", "user"),
+ "default", "default"),
+ &pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId6,
+ AppId: "default",
+ ServiceName:
"query_instance_shared_provider_ms",
+ VersionRule: "1.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Instances))
+ assert.Equal(t, instanceId5, respFind.Instances[0].InstanceId)
+
+ respFind, err = ms.MicroService().FindInstances(getContext(),
&pb.FindInstancesRequest{
+ ConsumerServiceId: serviceId7,
+ AppId: "default",
+ ServiceName: "query_instance_shared_provider_ms",
+ VersionRule: "1.0.0",
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Instances))
+ assert.Equal(t, instanceId5, respFind.Instances[0].InstanceId)
+
+ log.Info("query same domain deps")
+ // todo finish ut after implementing GetConsumerDependencies
interface
+
+ core.Service.Environment = pb.ENV_DEV
+ })
+
+ t.Run("batch query instances", func(t *testing.T) {
+ log.Info("find with version rule")
+ respFind, err := ms.MicroService().BatchFind(getContext(),
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId1,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_service_ms",
+ Version: "latest",
+ },
+ },
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_service_ms",
+ Version: "1.0.0+",
+ },
+ },
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_service_ms",
+ Version: "0.0.0",
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, int64(0), respFind.Services.Updated[0].Index)
+ assert.Equal(t, instanceId2,
respFind.Services.Updated[0].Instances[0].InstanceId)
+ assert.Equal(t, int64(1), respFind.Services.Updated[1].Index)
+ assert.Equal(t, instanceId2,
respFind.Services.Updated[1].Instances[0].InstanceId)
+ assert.Equal(t, int64(2),
respFind.Services.Failed[0].Indexes[0])
+ assert.Equal(t, scerr.ErrServiceNotExists,
respFind.Services.Failed[0].Error.Code)
+
+ log.Info("find with env")
+ respFind, err = ms.MicroService().BatchFind(getContext(),
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId4,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_diff_env_service_ms",
+ Version: "1.0.0",
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Services.Updated[0].Instances))
+ assert.Equal(t, instanceId4,
respFind.Services.Updated[0].Instances[0].InstanceId)
+
+ log.Info("find with rev")
+ ctx := util.SetContext(getContext(), util.CtxNocache, "")
+ respFind, err = ms.MicroService().BatchFind(ctx,
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId8,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_with_rev_ms",
+ Version: "1.0.0",
+ },
+ },
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"batch_query_instance_with_rev_ms",
+ Version: "1.0.0",
+ },
+ },
+ },
+ Instances: []*pb.FindInstance{
+ {
+ Instance: &pb.HeartbeatSetElement{
+ ServiceId: serviceId9,
+ InstanceId: instanceId9,
+ },
+ },
+ {
+ Instance: &pb.HeartbeatSetElement{
+ ServiceId: serviceId8,
+ InstanceId: instanceId8,
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ rev := respFind.Services.Updated[0].Rev
+ assert.Equal(t, int64(0), respFind.Services.Updated[0].Index)
+ assert.Equal(t, int64(1), respFind.Services.Updated[1].Index)
+ assert.Equal(t, instanceId8,
respFind.Services.Updated[0].Instances[0].InstanceId)
+ assert.Equal(t, instanceId9,
respFind.Services.Updated[1].Instances[0].InstanceId)
+ assert.NotEqual(t, 0, len(rev))
+ instanceRev := respFind.Instances.Updated[0].Rev
+ assert.Equal(t, int64(0), respFind.Instances.Updated[0].Index)
+ assert.Equal(t, int64(1), respFind.Instances.Updated[1].Index)
+ assert.Equal(t, instanceId9,
respFind.Instances.Updated[0].Instances[0].InstanceId)
+ assert.Equal(t, instanceId8,
respFind.Instances.Updated[1].Instances[0].InstanceId)
+ assert.NotEqual(t, 0, len(instanceRev))
+
+ respFind, err = ms.MicroService().BatchFind(ctx,
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId8,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_with_rev_ms",
+ Version: "1.0.0",
+ },
+ Rev: "x",
+ },
+ },
+ Instances: []*pb.FindInstance{
+ {
+ Instance: &pb.HeartbeatSetElement{
+ ServiceId: serviceId9,
+ InstanceId: instanceId9,
+ },
+ Rev: "x",
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, instanceId8,
respFind.Services.Updated[0].Instances[0].InstanceId)
+ assert.Equal(t, respFind.Services.Updated[0].Rev, rev)
+ assert.Equal(t, instanceId9,
respFind.Instances.Updated[0].Instances[0].InstanceId)
+ assert.Equal(t, instanceRev, respFind.Instances.Updated[0].Rev)
+
+ respFind, err = ms.MicroService().BatchFind(ctx,
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId8,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_with_rev_ms",
+ Version: "1.0.0",
+ },
+ Rev: rev,
+ },
+ },
+ Instances: []*pb.FindInstance{
+ {
+ Instance: &pb.HeartbeatSetElement{
+ ServiceId: serviceId9,
+ InstanceId: instanceId9,
+ },
+ Rev: instanceRev,
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, int64(0), respFind.Services.NotModified[0])
+ assert.Equal(t, int64(0), respFind.Instances.NotModified[0])
+
+ log.Info("find should return 200 even if consumer is diff apps")
+ respFind, err = ms.MicroService().BatchFind(getContext(),
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId3,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId:
"query_instance_ms",
+ ServiceName:
"query_instance_service_ms",
+ Version: "1.0.5",
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 0, len(respFind.Services.Updated[0].Instances))
+
+ log.Info("shared service discovery")
+ _ = os.Setenv("CSE_SHARED_SERVICES",
"query_instance_shared_provider_ms")
+ core.SetSharedMode()
+ core.Service.Environment = pb.ENV_PROD
+ respFind, err = ms.MicroService().BatchFind(
+ util.SetTargetDomainProject(
+
util.SetDomainProject(util.CloneContext(getContext()), "user", "user"),
+ "default", "default"),
+ &pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId6,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId: "default",
+ ServiceName:
"query_instance_shared_provider_ms",
+ Version: "1.0.0",
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Services.Updated[0].Instances))
+ assert.Equal(t, instanceId5,
respFind.Services.Updated[0].Instances[0].InstanceId)
+
+ respFind, err = ms.MicroService().BatchFind(getContext(),
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId7,
+ Services: []*pb.FindService{
+ {
+ Service: &pb.MicroServiceKey{
+ AppId: "default",
+ ServiceName:
"query_instance_shared_provider_ms",
+ Version: "1.0.0",
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Services.Updated[0].Instances))
+ assert.Equal(t, instanceId5,
respFind.Services.Updated[0].Instances[0].InstanceId)
+
+ respFind, err =
ms.MicroService().BatchFind(util.SetTargetDomainProject(
+ util.SetDomainProject(util.CloneContext(getContext()),
"user", "user"),
+ "default", "default"),
+ &pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId6,
+ Instances: []*pb.FindInstance{
+ {
+ Instance:
&pb.HeartbeatSetElement{
+ ServiceId: serviceId5,
+ InstanceId: instanceId5,
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, scerr.ErrServiceNotExists,
respFind.Instances.Failed[0].Error.Code)
+
+ respFind, err = ms.MicroService().BatchFind(getContext(),
&pb.BatchFindInstancesRequest{
+ ConsumerServiceId: serviceId7,
+ Instances: []*pb.FindInstance{
+ {
+ Instance: &pb.HeartbeatSetElement{
+ ServiceId: serviceId5,
+ InstanceId: instanceId5,
+ },
+ },
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ assert.Equal(t, 1, len(respFind.Instances.Updated[0].Instances))
+ assert.Equal(t, instanceId5,
respFind.Instances.Updated[0].Instances[0].InstanceId)
+
+ core.Service.Environment = pb.ENV_DEV
+ })
+
+ t.Run("query instances between diff dimensions", func(t *testing.T) {
+ log.Info("diff appId")
+ UTFunc := func(consumerId string, code int32) {
+ respFind, err :=
ms.MicroService().GetInstances(getContext(), &pb.GetInstancesRequest{
+ ConsumerServiceId: consumerId,
+ ProviderServiceId: serviceId2,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, code, respFind.Response.GetCode())
+ }
+
+ UTFunc(serviceId3, scerr.ErrServiceNotExists)
+
+ UTFunc(serviceId1, proto.Response_SUCCESS)
+
+ log.Info("diff env")
+ respFind, err := ms.MicroService().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: serviceId4,
+ ProviderServiceId: serviceId2,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respFind.Response.GetCode())
+ })
+}
+
+func TestInstance_GetOne(t *testing.T) {
+ ms.Install("etcd", func(opts ms.Options) (ms.DataSource, error) {
+ return etcd.NewDataSource(opts), nil
+ })
+ err := ms.Init(ms.Options{
+ Endpoint: "",
+ PluginImplName:
ms.ImplName(archaius.GetString("servicecomb.ms.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ var (
+ serviceId1 string
+ serviceId2 string
+ serviceId3 string
+ instanceId2 string
+ )
+
+ t.Run("register service and instances", func(t *testing.T) {
+ respCreateService, err :=
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "get_instance_ms",
+ ServiceName: "get_instance_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 =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "get_instance_ms",
+ ServiceName: "get_instance_service_ms",
+ Version: "1.0.5",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId2 = respCreateService.ServiceId
+
+ respCreateInstance, err :=
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId2,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "get:127.0.0.2:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ instanceId2 = respCreateInstance.InstanceId
+
+ respCreateService, err =
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "get_instance_cross_ms",
+ ServiceName: "get_instance_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())
+ serviceId3 = respCreateService.ServiceId
+ })
+
+ t.Run("get one instance when invalid request", func(t *testing.T) {
+ log.Info("find service itself")
+ resp, err := ms.MicroService().GetInstance(getContext(),
&pb.GetOneInstanceRequest{
+ ConsumerServiceId: serviceId2,
+ ProviderServiceId: serviceId2,
+ ProviderInstanceId: instanceId2,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+
+ log.Info("consumer does not exist")
+ resp, err = ms.MicroService().GetInstance(getContext(),
&pb.GetOneInstanceRequest{
+ ConsumerServiceId: "not-exist-id-ms",
+ ProviderServiceId: serviceId2,
+ ProviderInstanceId: instanceId2,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+ })
+
+ t.Run("get between diff apps", func(t *testing.T) {
+ resp, err := ms.MicroService().GetInstance(getContext(),
&pb.GetOneInstanceRequest{
+ ConsumerServiceId: serviceId3,
+ ProviderServiceId: serviceId2,
+ ProviderInstanceId: instanceId2,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, scerr.ErrInstanceNotExists,
resp.Response.GetCode())
+
+ respAll, err := ms.MicroService().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: serviceId3,
+ ProviderServiceId: serviceId2,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
respAll.Response.GetCode())
+ })
+
+ t.Run("get instances when request is invalid", func(t *testing.T) {
+ log.Info("consumer does not exist")
+ resp, err := ms.MicroService().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: "not-exist-service-ms",
+ ProviderServiceId: serviceId2,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+
+ log.Info("consumer does not exist")
+ resp, err = ms.MicroService().GetInstances(getContext(),
&pb.GetInstancesRequest{
+ ConsumerServiceId: serviceId1,
+ ProviderServiceId: serviceId2,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+}
+
+func TestInstance_Unregister(t *testing.T) {
+ ms.Install("etcd", func(opts ms.Options) (ms.DataSource, error) {
+ return etcd.NewDataSource(opts), nil
+ })
+ err := ms.Init(ms.Options{
+ Endpoint: "",
+ PluginImplName:
ms.ImplName(archaius.GetString("servicecomb.ms.name", "etcd")),
+ })
+ assert.NoError(t, err)
+
+ var (
+ serviceId string
+ instanceId string
+ )
+
+ t.Run("register service and instances", func(t *testing.T) {
+ respCreateService, err :=
ms.MicroService().RegisterService(getContext(), &pb.CreateServiceRequest{
+ Service: &pb.MicroService{
+ AppId: "unregister_instance_ms",
+ ServiceName: "unregister_instance_service_ms",
+ Version: "1.0.5",
+ Level: "FRONT",
+ Status: pb.MS_UP,
+ },
+ Tags: map[string]string{
+ "test": "test",
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateService.Response.GetCode())
+ serviceId = respCreateService.ServiceId
+
+ respCreateInstance, err :=
ms.MicroService().RegisterInstance(getContext(), &pb.RegisterInstanceRequest{
+ Instance: &pb.MicroServiceInstance{
+ ServiceId: serviceId,
+ HostName: "UT-HOST-MS",
+ Endpoints: []string{
+ "unregister:127.0.0.2:8080",
+ },
+ Status: pb.MSI_UP,
+ },
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS,
respCreateInstance.Response.GetCode())
+ instanceId = respCreateInstance.InstanceId
+ })
+
+ t.Run("unregister instance", func(t *testing.T) {
+ resp, err := ms.MicroService().UnregisterInstance(getContext(),
&pb.UnregisterInstanceRequest{
+ ServiceId: serviceId,
+ InstanceId: instanceId,
+ })
+ assert.NoError(t, err)
+ assert.Equal(t, proto.Response_SUCCESS, resp.Response.GetCode())
+ })
+
+ t.Run("unregister instance when request is invalid", func(t *testing.T)
{
+ log.Info("service id does not exist")
+ resp, err := ms.MicroService().UnregisterInstance(getContext(),
&pb.UnregisterInstanceRequest{
+ ServiceId: "not-exist-id-ms",
+ InstanceId: instanceId,
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+
+ log.Info("instance id does not exist")
+ resp, err = ms.MicroService().UnregisterInstance(getContext(),
&pb.UnregisterInstanceRequest{
+ ServiceId: serviceId,
+ InstanceId: "not-exist-id-ms",
+ })
+ assert.NoError(t, err)
+ assert.NotEqual(t, proto.Response_SUCCESS,
resp.Response.GetCode())
+ })
+}
+
func TestSchema_Create(t *testing.T) {
ms.Install("etcd", func(opts ms.Options) (ms.DataSource, error) {
return etcd.NewDataSource(opts), nil
diff --git a/server/service/ms/etcd/util.go b/server/service/ms/etcd/util.go
index 83736b7..d19fedd 100644
--- a/server/service/ms/etcd/util.go
+++ b/server/service/ms/etcd/util.go
@@ -18,6 +18,7 @@ package etcd
import (
"context"
"encoding/json"
+ errorsEx "github.com/apache/servicecomb-service-center/pkg/errors"
"github.com/apache/servicecomb-service-center/pkg/log"
pb "github.com/apache/servicecomb-service-center/pkg/registry"
"github.com/apache/servicecomb-service-center/pkg/util"
@@ -354,3 +355,38 @@ func preProcessRegisterInstance(ctx context.Context,
instance *pb.MicroServiceIn
instance.Version = microservice.Version
return nil
}
+
+func getHeartbeatFunc(ctx context.Context, domainProject string,
instancesHbRst chan<- *pb.InstanceHbRst, element *pb.HeartbeatSetElement)
func(context.Context) {
+ return func(_ context.Context) {
+ hbRst := &pb.InstanceHbRst{
+ ServiceId: element.ServiceId,
+ InstanceId: element.InstanceId,
+ ErrMessage: "",
+ }
+ _, _, err := serviceUtil.HeartbeatUtil(ctx, domainProject,
element.ServiceId, element.InstanceId)
+ if err != nil {
+ hbRst.ErrMessage = err.Error()
+ log.Errorf(err, "heartbeat set failed, %s/%s",
element.ServiceId, element.InstanceId)
+ }
+ instancesHbRst <- hbRst
+ }
+}
+
+func revokeInstance(ctx context.Context, domainProject string, serviceID
string, instanceID string) *scerr.Error {
+ leaseID, err := serviceUtil.GetLeaseID(ctx, domainProject, serviceID,
instanceID)
+ if err != nil {
+ return scerr.NewError(scerr.ErrUnavailableBackend, err.Error())
+ }
+ if leaseID == -1 {
+ return scerr.NewError(scerr.ErrInstanceNotExists, "Instance's
leaseId not exist.")
+ }
+
+ err = backend.Registry().LeaseRevoke(ctx, leaseID)
+ if err != nil {
+ if _, ok := err.(errorsEx.InternalError); !ok {
+ return scerr.NewError(scerr.ErrInstanceNotExists,
err.Error())
+ }
+ return scerr.NewError(scerr.ErrUnavailableBackend, err.Error())
+ }
+ return nil
+}