From a14207e67a7ac4676a06d0459858ff49e04769be Mon Sep 17 00:00:00 2001 From: peng-yongsheng <8082209@qq.com> Date: Tue, 13 Feb 2018 16:12:19 +0800 Subject: [PATCH] Instance heart beat tested with elastic search. OK. --- .../InstanceDiscoveryServiceHandler.java | 2 ++ .../provider/handler/mock/RegisterMock.java | 23 +++++++++++++++++++ .../handler/mock/TraceSegmentMock.java | 2 ++ 3 files changed, 27 insertions(+) diff --git a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/InstanceDiscoveryServiceHandler.java b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/InstanceDiscoveryServiceHandler.java index cde3e4c69..2b20ce391 100644 --- a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/InstanceDiscoveryServiceHandler.java +++ b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/main/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/InstanceDiscoveryServiceHandler.java @@ -66,6 +66,8 @@ public class InstanceDiscoveryServiceHandler extends InstanceDiscoveryServiceGrp int instanceId = request.getApplicationInstanceId(); long heartBeatTime = request.getHeartbeatTime(); this.instanceHeartBeatService.heartBeat(instanceId, heartBeatTime); + responseObserver.onNext(Downstream.getDefaultInstance()); + responseObserver.onCompleted(); } private String buildOsInfo(OSInfo osinfo) { diff --git a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/RegisterMock.java b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/RegisterMock.java index c27301eac..11c7f070a 100644 --- a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/RegisterMock.java +++ b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/RegisterMock.java @@ -20,8 +20,11 @@ package org.apache.skywalking.apm.collector.agent.grpc.provider.handler.mock; import io.grpc.ManagedChannel; import java.util.UUID; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import org.apache.skywalking.apm.network.proto.Application; import org.apache.skywalking.apm.network.proto.ApplicationInstance; +import org.apache.skywalking.apm.network.proto.ApplicationInstanceHeartbeat; import org.apache.skywalking.apm.network.proto.ApplicationInstanceMapping; import org.apache.skywalking.apm.network.proto.ApplicationMapping; import org.apache.skywalking.apm.network.proto.ApplicationRegisterServiceGrpc; @@ -31,6 +34,7 @@ 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.apache.skywalking.apm.network.proto.ServiceNameMappingCollection; +import org.apache.skywalking.apm.util.RunnableWithExceptionProtection; import org.joda.time.DateTime; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -94,6 +98,8 @@ class RegisterMock { serviceNameCollection.addElements(serviceNameElement); registerServiceName(serviceNameCollection); + + heartBeatScheduled(instanceMapping.getApplicationInstanceId()); } private void registerProvider() throws InterruptedException { @@ -136,6 +142,8 @@ class RegisterMock { serviceNameCollection.addElements(serviceNameElement); registerServiceName(serviceNameCollection); + + heartBeatScheduled(instanceMapping.getApplicationInstanceId()); } private void registerServiceName(ServiceNameCollection.Builder serviceNameCollection) throws InterruptedException { @@ -150,4 +158,19 @@ class RegisterMock { } while (serviceNameMappingCollection.getElementsCount() == 0 || serviceNameMappingCollection.getElements(0).getServiceId() == 0); } + + private void heartBeatScheduled(int instanceId) { + Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate( + new RunnableWithExceptionProtection(() -> heartBeat(instanceId), + t -> logger.error("instance heart beat scheduled error.", t)), 4, 1, TimeUnit.SECONDS); + } + + private void heartBeat(int instanceId) { + long now = System.currentTimeMillis(); + logger.debug("instance heart beat, instance id: {}, time: {}", instanceId, now); + ApplicationInstanceHeartbeat.Builder heartbeat = ApplicationInstanceHeartbeat.newBuilder(); + heartbeat.setApplicationInstanceId(instanceId); + heartbeat.setHeartbeatTime(now); + instanceDiscoveryServiceBlockingStub.heartbeat(heartbeat.build()); + } } diff --git a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/TraceSegmentMock.java b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/TraceSegmentMock.java index 8fe9ee82a..c290bd30a 100644 --- a/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/TraceSegmentMock.java +++ b/apm-collector/apm-collector-agent/agent-grpc/agent-grpc-provider/src/test/java/org/apache/skywalking/apm/collector/agent/grpc/provider/handler/mock/TraceSegmentMock.java @@ -83,6 +83,8 @@ public class TraceSegmentMock { while (sleeping.getValue()) { Thread.sleep(200); } + + Thread.sleep(200000); } static class Sleeping {