Merge pull request #816 from peng-yongsheng/feature/register_test

Application, Instance, Service register tested with elastic search.
This commit is contained in:
吴晟 Wu Sheng 2018-02-13 15:24:36 +08:00 committed by GitHub
commit abb3515963
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
6 changed files with 129 additions and 7 deletions

View File

@ -0,0 +1,104 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* 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 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
* 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
* limitations under the License.
*
*/
package org.apache.skywalking.apm.collector.agent.grpc.provider.handler.mock;
import io.grpc.ManagedChannel;
import java.util.UUID;
import org.apache.skywalking.apm.network.proto.Application;
import org.apache.skywalking.apm.network.proto.ApplicationInstance;
import org.apache.skywalking.apm.network.proto.ApplicationMapping;
import org.apache.skywalking.apm.network.proto.ApplicationRegisterServiceGrpc;
import org.apache.skywalking.apm.network.proto.InstanceDiscoveryServiceGrpc;
import org.apache.skywalking.apm.network.proto.OSInfo;
import org.apache.skywalking.apm.network.proto.ServiceNameCollection;
import org.apache.skywalking.apm.network.proto.ServiceNameDiscoveryServiceGrpc;
import org.apache.skywalking.apm.network.proto.ServiceNameElement;
import org.joda.time.DateTime;
/**
* @author peng-yongsheng
*/
class RegisterMock {
private ApplicationRegisterServiceGrpc.ApplicationRegisterServiceBlockingStub applicationRegisterServiceBlockingStub;
private InstanceDiscoveryServiceGrpc.InstanceDiscoveryServiceBlockingStub instanceDiscoveryServiceBlockingStub;
private ServiceNameDiscoveryServiceGrpc.ServiceNameDiscoveryServiceBlockingStub serviceNameDiscoveryServiceBlockingStub;
void mock(ManagedChannel channel) {
applicationRegisterServiceBlockingStub = ApplicationRegisterServiceGrpc.newBlockingStub(channel);
instanceDiscoveryServiceBlockingStub = InstanceDiscoveryServiceGrpc.newBlockingStub(channel);
serviceNameDiscoveryServiceBlockingStub = ServiceNameDiscoveryServiceGrpc.newBlockingStub(channel);
registerConsumer();
registerProvider();
}
private void registerConsumer() {
Application.Builder application = Application.newBuilder();
application.setApplicationCode("dubbox-consumer");
ApplicationMapping applicationMapping = applicationRegisterServiceBlockingStub.applicationCodeRegister(application.build());
ApplicationInstance.Builder instance = ApplicationInstance.newBuilder();
instance.setApplicationId(applicationMapping.getApplication().getValue());
instance.setAgentUUID(UUID.randomUUID().toString());
instance.setRegisterTime(new DateTime("2017-01-01T00:01:01.001").getMillis());
OSInfo.Builder osInfo = OSInfo.newBuilder();
osInfo.setHostname("pengys");
osInfo.setOsName("MacOS XX");
osInfo.setProcessNo(1001);
osInfo.addIpv4S("10.0.0.3");
osInfo.addIpv4S("10.0.0.4");
instance.setOsinfo(osInfo);
instanceDiscoveryServiceBlockingStub.registerInstance(instance.build());
ServiceNameCollection.Builder serviceNameCollection = ServiceNameCollection.newBuilder();
ServiceNameElement.Builder serviceNameElement = ServiceNameElement.newBuilder();
serviceNameElement.setApplicationId(applicationMapping.getApplication().getValue());
serviceNameElement.setServiceName("org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()");
serviceNameCollection.addElements(serviceNameElement);
serviceNameDiscoveryServiceBlockingStub.discovery(serviceNameCollection.build());
}
private void registerProvider() {
Application.Builder application = Application.newBuilder();
application.setApplicationCode("dubbox-provider");
ApplicationMapping applicationMapping = applicationRegisterServiceBlockingStub.applicationCodeRegister(application.build());
ApplicationInstance.Builder instance = ApplicationInstance.newBuilder();
instance.setApplicationId(applicationMapping.getApplication().getValue());
instance.setAgentUUID(UUID.randomUUID().toString());
instance.setRegisterTime(new DateTime("2017-01-01T00:01:01.001").getMillis());
OSInfo.Builder osInfo = OSInfo.newBuilder();
osInfo.setHostname("peng-yongsheng");
osInfo.setOsName("MacOS X");
osInfo.setProcessNo(1000);
osInfo.addIpv4S("10.0.0.1");
osInfo.addIpv4S("10.0.0.2");
instance.setOsinfo(osInfo);
instanceDiscoveryServiceBlockingStub.registerInstance(instance.build());
ServiceNameCollection.Builder serviceNameCollection = ServiceNameCollection.newBuilder();
ServiceNameElement.Builder serviceNameElement = ServiceNameElement.newBuilder();
serviceNameElement.setApplicationId(applicationMapping.getApplication().getValue());
serviceNameElement.setServiceName("org.skywaking.apm.testcase.dubbo.services.GreetService.doBusiness()");
serviceNameCollection.addElements(serviceNameElement);
serviceNameDiscoveryServiceBlockingStub.discovery(serviceNameCollection.build());
}
}

View File

@ -36,9 +36,13 @@ public class TraceSegmentMock {
private static final Logger logger = LoggerFactory.getLogger(TraceSegmentMock.class);
public static void main(String[] args) throws InterruptedException {
ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).usePlaintext(true).build();
RegisterMock registerMock = new RegisterMock();
registerMock.mock(channel);
Sleeping sleeping = new Sleeping();
ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 11800).usePlaintext(true).build();
TraceSegmentServiceGrpc.TraceSegmentServiceStub stub = TraceSegmentServiceGrpc.newStub(channel);
StreamObserver<UpstreamSegment> segmentStreamObserver = stub.collect(new StreamObserver<Downstream>() {
@Override public void onNext(Downstream downstream) {

View File

@ -55,7 +55,14 @@ public class ApplicationRegisterSerialWorker extends AbstractLocalAsyncWorker<Ap
@Override protected void onWork(Application application) throws WorkerException {
logger.debug("register application, application code: {}", application.getApplicationCode());
int applicationId = applicationCacheService.getApplicationIdByCode(application.getApplicationCode());
int applicationId;
if (BooleanUtils.valueToBoolean(application.getIsAddress())) {
applicationId = applicationCacheService.getApplicationIdByAddressId(application.getAddressId());
} else {
applicationId = applicationCacheService.getApplicationIdByCode(application.getApplicationCode());
}
if (applicationId == 0) {
Application newApplication;

View File

@ -55,7 +55,14 @@ public class InstanceRegisterSerialWorker extends AbstractLocalAsyncWorker<Insta
@Override protected void onWork(Instance instance) throws WorkerException {
logger.debug("register instance, application id: {}, agentUUID: {}", instance.getApplicationId(), instance.getAgentUUID());
int instanceId = instanceCacheService.getInstanceIdByAgentUUID(instance.getApplicationId(), instance.getAgentUUID());
int instanceId;
if (BooleanUtils.valueToBoolean(instance.getIsAddress())) {
instanceId = instanceCacheService.getInstanceIdByAddressId(instance.getApplicationId(), instance.getAddressId());
} else {
instanceId = instanceCacheService.getInstanceIdByAgentUUID(instance.getApplicationId(), instance.getAgentUUID());
}
if (instanceId == 0) {
Instance newInstance;

View File

@ -50,7 +50,7 @@ public class ApplicationEsCacheDAO extends EsDAO implements IApplicationCacheDAO
ElasticSearchClient client = getClient();
SearchRequestBuilder searchRequestBuilder = client.prepareSearch(ApplicationTable.TABLE);
searchRequestBuilder.setTypes("type");
searchRequestBuilder.setTypes(ApplicationTable.TABLE_TYPE);
searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH);
BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery();
@ -88,7 +88,7 @@ public class ApplicationEsCacheDAO extends EsDAO implements IApplicationCacheDAO
ElasticSearchClient client = getClient();
SearchRequestBuilder searchRequestBuilder = client.prepareSearch(ApplicationTable.TABLE);
searchRequestBuilder.setTypes("type");
searchRequestBuilder.setTypes(ApplicationTable.TABLE_TYPE);
searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH);
BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery();

View File

@ -53,7 +53,7 @@ public class InstanceEsCacheDAO extends EsDAO implements IInstanceCacheDAO {
ElasticSearchClient client = getClient();
SearchRequestBuilder searchRequestBuilder = client.prepareSearch(InstanceTable.TABLE);
searchRequestBuilder.setTypes("type");
searchRequestBuilder.setTypes(InstanceTable.TABLE_TYPE);
searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH);
BoolQueryBuilder builder = QueryBuilders.boolQuery();
builder.must().add(QueryBuilders.termQuery(InstanceTable.COLUMN_APPLICATION_ID, applicationId));
@ -74,7 +74,7 @@ public class InstanceEsCacheDAO extends EsDAO implements IInstanceCacheDAO {
ElasticSearchClient client = getClient();
SearchRequestBuilder searchRequestBuilder = client.prepareSearch(InstanceTable.TABLE);
searchRequestBuilder.setTypes("type");
searchRequestBuilder.setTypes(InstanceTable.TABLE_TYPE);
searchRequestBuilder.setSearchType(SearchType.QUERY_THEN_FETCH);
BoolQueryBuilder builder = QueryBuilders.boolQuery();
builder.must().add(QueryBuilders.termQuery(InstanceTable.COLUMN_APPLICATION_ID, applicationId));