From 5d5100ef1bfdee3e484419289dcf6381fb284710 Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Thu, 1 Nov 2018 18:02:09 +0200 Subject: [PATCH] Improved logging --- .../server/actors/app/AppActor.java | 1 - .../server/actors/device/DeviceActor.java | 2 -- .../device/DeviceActorMessageProcessor.java | 3 ++ .../actors/rpc/BasicRpcSessionListener.java | 10 +++--- .../server/actors/rpc/RpcManagerActor.java | 1 - .../RuleChainActorMessageProcessor.java | 30 +++++++++++++--- .../actors/ruleChain/RuleNodeActor.java | 10 +++--- .../RuleNodeActorMessageProcessor.java | 9 +++-- .../server/actors/service/ComponentActor.java | 8 +++-- .../actors/service/ContextAwareActor.java | 9 ++--- .../actors/shared/ComponentMsgProcessor.java | 2 ++ .../server/actors/tenant/TenantActor.java | 1 - .../service/script/RemoteJsInvokeService.java | 4 +-- .../RemoteRuleEngineTransportService.java | 12 ++++--- .../server/kafka/TBKafkaProducerTemplate.java | 5 +-- .../server/kafka/TbKafkaRequestTemplate.java | 20 +++++++++-- msa/black-box-tests/README.md | 4 ++- .../server/msa/AbstractContainerTest.java | 31 +++++++++++++++++ .../org/thingsboard/server/msa/WsClient.java | 34 ++++++++++++++----- .../msa/connectivity/MqttClientTest.java | 12 +++++-- 20 files changed, 157 insertions(+), 51 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java index 15e9e9d175..6df37d6731 100644 --- a/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/app/AppActor.java @@ -51,7 +51,6 @@ import java.util.HashMap; import java.util.Map; import java.util.Optional; -@Slf4j public class AppActor extends RuleChainManagerActor { private static final TenantId SYSTEM_TENANT = new TenantId(ModelConstants.NULL_UUID); diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java index f450027fb6..f53a410d5f 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActor.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.actors.device; -import lombok.extern.slf4j.Slf4j; import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; import org.thingsboard.rule.engine.api.msg.DeviceNameOrTypeUpdateMsg; import org.thingsboard.server.actors.ActorSystemContext; @@ -29,7 +28,6 @@ import org.thingsboard.server.service.rpc.ToDeviceRpcRequestActorMsg; import org.thingsboard.server.service.rpc.ToServerRpcResponseActorMsg; import org.thingsboard.server.service.transport.msg.TransportToDeviceActorMsgWrapper; -@Slf4j public class DeviceActor extends ContextAwareActor { private final DeviceActorMessageProcessor processor; diff --git a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java index 47417d1fd7..41d1ea08c6 100644 --- a/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java @@ -348,9 +348,12 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor { int requestId = msg.getMsg().getRequestId(); ToServerRpcRequestMetadata data = toServerRpcPendingMap.remove(requestId); if (data != null) { + log.debug("[{}] Pushing reply to [{}][{}]!", deviceId, data.getNodeId(), data.getSessionId()); sendToTransport(TransportProtos.ToServerRpcResponseMsg.newBuilder() .setRequestId(requestId).setPayload(msg.getMsg().getData()).build() , data.getSessionId(), data.getNodeId()); + } else { + log.debug("[{}][{}] Pending RPC request to server not found!", deviceId, requestId); } } diff --git a/application/src/main/java/org/thingsboard/server/actors/rpc/BasicRpcSessionListener.java b/application/src/main/java/org/thingsboard/server/actors/rpc/BasicRpcSessionListener.java index 1eba066960..760d4a6d96 100644 --- a/application/src/main/java/org/thingsboard/server/actors/rpc/BasicRpcSessionListener.java +++ b/application/src/main/java/org/thingsboard/server/actors/rpc/BasicRpcSessionListener.java @@ -40,7 +40,7 @@ public class BasicRpcSessionListener implements GrpcSessionListener { @Override public void onConnected(GrpcSession session) { - log.info("{} session started -> {}", getType(session), session.getRemoteServer()); + log.info("[{}][{}] session started", session.getRemoteServer(), getType(session)); if (!session.isClient()) { manager.tell(new RpcSessionConnectedMsg(session.getRemoteServer(), session.getSessionId()), self); } @@ -48,21 +48,19 @@ public class BasicRpcSessionListener implements GrpcSessionListener { @Override public void onDisconnected(GrpcSession session) { - log.info("{} session closed -> {}", getType(session), session.getRemoteServer()); + log.info("[{}][{}] session closed", session.getRemoteServer(), getType(session)); manager.tell(new RpcSessionDisconnectedMsg(session.isClient(), session.getRemoteServer()), self); } @Override public void onReceiveClusterGrpcMsg(GrpcSession session, ClusterAPIProtos.ClusterMessage clusterMessage) { - log.trace("{} Service [{}] received session actor msg {}", getType(session), - session.getRemoteServer(), - clusterMessage); + log.trace("Received session actor msg from [{}][{}]: {}", session.getRemoteServer(), getType(session), clusterMessage); service.onReceivedMsg(session.getRemoteServer(), clusterMessage); } @Override public void onError(GrpcSession session, Throwable t) { - log.warn("{} session got error -> {}", getType(session), session.getRemoteServer(), t); + log.warn("[{}][{}] session got error -> {}", session.getRemoteServer(), getType(session), t); manager.tell(new RpcSessionClosedMsg(session.isClient(), session.getRemoteServer()), self); session.close(); } diff --git a/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java b/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java index 922d2ba988..9fa087e630 100644 --- a/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/rpc/RpcManagerActor.java @@ -36,7 +36,6 @@ import java.util.*; /** * @author Andrew Shvayka */ -@Slf4j public class RpcManagerActor extends ContextAwareActor { private final Map sessionActors; diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java index e0f7e4ad26..e4e82c55e5 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/RuleChainActorMessageProcessor.java @@ -69,6 +69,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor(); this.nodeRoutes = new HashMap<>(); this.service = systemContext.getRuleChainService(); + this.ruleChainName = ruleChainId.toString(); } @Override - public void start(ActorContext context) throws Exception { + public String getComponentName() { + return null; + } + + @Override + public void start(ActorContext context) { if (!started) { RuleChain ruleChain = service.findRuleChainById(entityId); + ruleChainName = ruleChain.getName(); List ruleNodeList = service.getRuleChainNodes(entityId); + log.trace("[{}][{}] Starting rule chain with {} nodes", tenantId, entityId, ruleNodeList.size()); // Creating and starting the actors; for (RuleNode ruleNode : ruleNodeList) { + log.trace("[{}][{}] Creating rule node [{}]: {}", tenantId, entityId, ruleNode.getName(), ruleNode); ActorRef ruleNodeActor = createRuleNodeActor(context, ruleNode); nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode)); } @@ -98,16 +108,19 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor ruleNodeList = service.getRuleChainNodes(entityId); - + log.trace("[{}][{}] Updating rule chain with {} nodes", tenantId, entityId, ruleNodeList.size()); for (RuleNode ruleNode : ruleNodeList) { RuleNodeCtx existing = nodeActors.get(ruleNode.getId()); if (existing == null) { + log.trace("[{}][{}] Creating rule node [{}]: {}", tenantId, entityId, ruleNode.getName(), ruleNode); ActorRef ruleNodeActor = createRuleNodeActor(context, ruleNode); nodeActors.put(ruleNode.getId(), new RuleNodeCtx(tenantId, self, ruleNodeActor, ruleNode)); } else { + log.trace("[{}][{}] Updating rule node [{}]: {}", tenantId, entityId, ruleNode.getName(), ruleNode); existing.setSelf(ruleNode); existing.getSelfActor().tell(new ComponentLifecycleMsg(tenantId, existing.getSelf().getId(), ComponentLifecycleEvent.UPDATED), self); } @@ -116,6 +129,7 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor existingNodes = ruleNodeList.stream().map(RuleNode::getId).collect(Collectors.toSet()); List removedRules = nodeActors.keySet().stream().filter(node -> !existingNodes.contains(node)).collect(Collectors.toList()); removedRules.forEach(ruleNodeId -> { + log.trace("[{}][{}] Removing rule node [{}]", tenantId, entityId, ruleNodeId); RuleNodeCtx removed = nodeActors.remove(ruleNodeId); removed.getSelfActor().tell(new ComponentLifecycleMsg(tenantId, removed.getSelf().getId(), ComponentLifecycleEvent.DELETED), self); }); @@ -124,7 +138,8 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor relations = service.getRuleNodeRelations(ruleNode.getId()); + log.trace("[{}][{}][{}] Processing rule node relations [{}]", tenantId, entityId, ruleNode.getId(), relations.size()); if (relations.size() == 0) { nodeRoutes.put(ruleNode.getId(), Collections.emptyList()); } else { for (EntityRelation relation : relations) { + log.trace("[{}][{}][{}] Processing rule node relation [{}]", tenantId, entityId, ruleNode.getId(), relation.getTo()); if (relation.getTo().getEntityType() == EntityType.RULE_NODE) { RuleNodeCtx ruleNodeCtx = nodeActors.get(new RuleNodeId(relation.getTo().getId())); if (ruleNodeCtx == null) { @@ -232,17 +249,20 @@ public class RuleChainActorMessageProcessor extends ComponentMsgProcessor { private final RuleChainId ruleChainId; @@ -62,7 +60,9 @@ public class RuleNodeActor extends ComponentActor componentClazz = Class.forName(ruleNode.getType()); TbNode tbNode = (TbNode) (componentClazz.newInstance()); diff --git a/application/src/main/java/org/thingsboard/server/actors/service/ComponentActor.java b/application/src/main/java/org/thingsboard/server/actors/service/ComponentActor.java index ff54fc8fa1..1f084e1bb7 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/ComponentActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/ComponentActor.java @@ -31,7 +31,6 @@ import org.thingsboard.server.common.msg.plugin.ComponentLifecycleMsg; /** * @author Andrew Shvayka */ -@Slf4j public abstract class ComponentActor> extends ContextAwareActor { private long lastPersistedErrorTs = 0L; @@ -54,6 +53,7 @@ public abstract class ComponentActor getErrorPersistFrequency()) { diff --git a/application/src/main/java/org/thingsboard/server/actors/service/ContextAwareActor.java b/application/src/main/java/org/thingsboard/server/actors/service/ContextAwareActor.java index c9f1dcde3e..1c7e2d2664 100644 --- a/application/src/main/java/org/thingsboard/server/actors/service/ContextAwareActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/service/ContextAwareActor.java @@ -17,15 +17,16 @@ package org.thingsboard.server.actors.service; import akka.actor.Terminated; import akka.actor.UntypedActor; -import akka.event.Logging; -import akka.event.LoggingAdapter; -import lombok.extern.slf4j.Slf4j; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.msg.TbActorMsg; -@Slf4j + public abstract class ContextAwareActor extends UntypedActor { + protected final Logger log = LoggerFactory.getLogger(getClass()); + public static final int ENTITY_PACK_LIMIT = 1024; protected final ActorSystemContext systemContext; diff --git a/application/src/main/java/org/thingsboard/server/actors/shared/ComponentMsgProcessor.java b/application/src/main/java/org/thingsboard/server/actors/shared/ComponentMsgProcessor.java index 9e5d063223..46a76e6d43 100644 --- a/application/src/main/java/org/thingsboard/server/actors/shared/ComponentMsgProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/shared/ComponentMsgProcessor.java @@ -44,6 +44,8 @@ public abstract class ComponentMsgProcessor extends Abstract this.entityId = id; } + public abstract String getComponentName(); + public abstract void start(ActorContext context) throws Exception; public abstract void stop(ActorContext context) throws Exception; diff --git a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java index b76e2984ae..6c65dbe531 100644 --- a/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/tenant/TenantActor.java @@ -48,7 +48,6 @@ import scala.concurrent.duration.Duration; import java.util.HashMap; import java.util.Map; -@Slf4j public class TenantActor extends RuleChainManagerActor { private final TenantId tenantId; diff --git a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java index db31bda942..957624c4da 100644 --- a/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java +++ b/application/src/main/java/org/thingsboard/server/service/script/RemoteJsInvokeService.java @@ -71,7 +71,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { private int maxErrors; private TbKafkaRequestTemplate kafkaTemplate; - protected Map scriptIdToBodysMap = new ConcurrentHashMap<>(); + private Map scriptIdToBodysMap = new ConcurrentHashMap<>(); @PostConstruct public void init() { @@ -100,7 +100,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { responseBuilder.settings(kafkaSettings); responseBuilder.topic(responseTopicPrefix + "." + nodeIdProvider.getNodeId()); responseBuilder.clientId("js-" + nodeIdProvider.getNodeId()); - responseBuilder.groupId("rule-engine-node"); + responseBuilder.groupId("rule-engine-node-" + nodeIdProvider.getNodeId()); responseBuilder.autoCommit(true); responseBuilder.autoCommitIntervalMs(autoCommitInterval); responseBuilder.decoder(new RemoteJsResponseDecoder()); diff --git a/application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java b/application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java index 6c903cc955..a8ab2cd625 100644 --- a/application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java +++ b/application/src/main/java/org/thingsboard/server/service/transport/RemoteRuleEngineTransportService.java @@ -149,6 +149,7 @@ public class RemoteRuleEngineTransportService implements RuleEngineTransportServ records.forEach(record -> { try { ToRuleEngineMsg toRuleEngineMsg = ruleEngineConsumer.decode(record); + log.trace("Forwarding message to rule engine {}", toRuleEngineMsg); if (toRuleEngineMsg.hasToDeviceActorMsg()) { forwardToDeviceActor(toRuleEngineMsg.getToDeviceActorMsg()); } @@ -175,18 +176,21 @@ public class RemoteRuleEngineTransportService implements RuleEngineTransportServ @Override public void process(String nodeId, DeviceActorToTransportMsg msg, Runnable onSuccess, Consumer onFailure) { - notificationsProducer.send(notificationsTopic + "." + nodeId, - new UUID(msg.getSessionIdMSB(), msg.getSessionIdLSB()).toString(), - ToTransportMsg.newBuilder().setToDeviceSessionMsg(msg).build() - , new QueueCallbackAdaptor(onSuccess, onFailure)); + String topic = notificationsTopic + "." + nodeId; + UUID sessionId = new UUID(msg.getSessionIdMSB(), msg.getSessionIdLSB()); + ToTransportMsg transportMsg = ToTransportMsg.newBuilder().setToDeviceSessionMsg(msg).build(); + log.trace("[{}][{}] Pushing session data to topic: {}", topic, sessionId, transportMsg); + notificationsProducer.send(topic, sessionId.toString(), transportMsg, new QueueCallbackAdaptor(onSuccess, onFailure)); } private void forwardToDeviceActor(TransportToDeviceActorMsg toDeviceActorMsg) { TransportToDeviceActorMsgWrapper wrapper = new TransportToDeviceActorMsgWrapper(toDeviceActorMsg); Optional address = routingService.resolveById(wrapper.getDeviceId()); if (address.isPresent()) { + log.trace("[{}] Pushing message to remote server: {}", address.get(), toDeviceActorMsg); rpcService.tell(encodingService.convertToProtoDataMessage(address.get(), wrapper)); } else { + log.trace("Pushing message to local server: {}", toDeviceActorMsg); actorContext.getAppActor().tell(wrapper, ActorRef.noSender()); } } diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java index bd42f310a6..ee652f455c 100644 --- a/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/kafka/TBKafkaProducerTemplate.java @@ -77,9 +77,10 @@ public class TBKafkaProducerTemplate { result.all().get(); } catch (Exception e) { if ((e instanceof TopicExistsException) || (e.getCause() != null && e.getCause() instanceof TopicExistsException)) { - log.trace("[{}] Topic already exists: ", defaultTopic); + log.trace("[{}] Topic already exists.", defaultTopic); } else { - log.trace("[{}] Failed to create topic: {}", defaultTopic, e.getMessage(), e); + log.info("[{}] Failed to create topic: {}", defaultTopic, e.getMessage(), e); + throw new RuntimeException(e); } } //Maybe this should not be cached, but we don't plan to change size of partitions diff --git a/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java index 77ad0331ab..86fda76624 100644 --- a/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/kafka/TbKafkaRequestTemplate.java @@ -23,6 +23,9 @@ import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.admin.CreateTopicsResult; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.producer.Callback; +import org.apache.kafka.clients.producer.RecordMetadata; +import org.apache.kafka.common.errors.TopicExistsException; import org.apache.kafka.common.header.Header; import org.apache.kafka.common.header.internals.RecordHeader; @@ -83,7 +86,13 @@ public class TbKafkaRequestTemplate extends AbstractTbKafkaTe CreateTopicsResult result = admin.createTopic(new NewTopic(responseTemplate.getTopic(), 1, (short) 1)); result.all().get(); } catch (Exception e) { - log.trace("Failed to create topic: {}", e.getMessage(), e); + if ((e instanceof TopicExistsException) || (e.getCause() != null && e.getCause() instanceof TopicExistsException)) { + log.trace("[{}] Topic already exists. ", responseTemplate.getTopic()); + } else { + log.info("[{}] Failed to create topic: {}", responseTemplate.getTopic(), e.getMessage(), e); + throw new RuntimeException(e); + } + } this.requestTemplate.init(); tickTs = System.currentTimeMillis(); @@ -93,6 +102,7 @@ public class TbKafkaRequestTemplate extends AbstractTbKafkaTe while (!stopped) { ConsumerRecords responses = responseTemplate.poll(Duration.ofMillis(pollInterval)); responses.forEach(response -> { + log.trace("Received response to Kafka Template request: {}", response); Header requestIdHeader = response.headers().lastHeader(TbKafkaSettings.REQUEST_ID_HEADER); Response decodedResponse = null; UUID requestId = null; @@ -160,7 +170,13 @@ public class TbKafkaRequestTemplate extends AbstractTbKafkaTe SettableFuture future = SettableFuture.create(); pendingRequests.putIfAbsent(requestId, new ResponseMetaData<>(tickTs + maxRequestTimeout, future)); request = requestTemplate.enrich(request, responseTemplate.getTopic(), requestId); - requestTemplate.send(key, request, headers, null); + requestTemplate.send(key, request, headers, (metadata, exception) -> { + if (exception != null) { + log.trace("[{}] Failed to post the request", requestId, exception); + } else { + log.trace("[{}] Posted the request", requestId, metadata); + } + }); return future; } diff --git a/msa/black-box-tests/README.md b/msa/black-box-tests/README.md index c26d9c5fc6..044fe00255 100644 --- a/msa/black-box-tests/README.md +++ b/msa/black-box-tests/README.md @@ -18,4 +18,6 @@ As result, in REPOSITORY column, next images should be present: - Run the black box tests in the [msa/black-box-tests](../black-box-tests) directory: - mvn clean install -DblackBoxTests.skip=false \ No newline at end of file + mvn clean install -DblackBoxTests.skip=false + + diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java index 5ebab780be..4c3ac8362c 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/AbstractContainerTest.java @@ -33,6 +33,9 @@ import org.apache.http.impl.conn.PoolingHttpClientConnectionManager; import org.apache.http.ssl.SSLContextBuilder; import org.apache.http.ssl.SSLContexts; import org.junit.*; +import org.junit.rules.TestRule; +import org.junit.rules.TestWatcher; +import org.junit.runner.Description; import org.springframework.http.client.HttpComponentsClientHttpRequestFactory; import org.thingsboard.client.tools.RestClient; import org.thingsboard.server.common.data.Device; @@ -60,6 +63,33 @@ public abstract class AbstractContainerTest { restClient.getRestTemplate().setRequestFactory(getRequestFactoryForSelfSignedCert()); } + @Rule + public TestRule watcher = new TestWatcher() { + protected void starting(Description description) { + log.info("================================================="); + log.info("STARTING TEST: {}" , description.getMethodName()); + log.info("================================================="); + } + + /** + * Invoked when a test succeeds + */ + protected void succeeded(Description description) { + log.info("================================================="); + log.info("SUCCEEDED TEST: {}" , description.getMethodName()); + log.info("================================================="); + } + + /** + * Invoked when a test fails + */ + protected void failed(Throwable e, Description description) { + log.info("================================================="); + log.info("FAILED TEST: {}" , description.getMethodName(), e); + log.info("================================================="); + } + }; + protected Device createDevice(String name) { return restClient.createDevice(name + RandomStringUtils.randomAlphanumeric(7), "DEFAULT"); } @@ -82,6 +112,7 @@ public abstract class AbstractContainerTest { JsonObject wsRequest = new JsonObject(); wsRequest.add(property.toString(), cmd); wsClient.send(wsRequest.toString()); + wsClient.waitForFirstReply(); return wsClient; } diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/WsClient.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/WsClient.java index a9835ed618..fa3a63a1e6 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/WsClient.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/WsClient.java @@ -31,9 +31,11 @@ public class WsClient extends WebSocketClient { private static final ObjectMapper mapper = new ObjectMapper(); private WsTelemetryResponse message; - private CountDownLatch latch = new CountDownLatch(1);; + private volatile boolean firstReplyReceived; + private CountDownLatch firstReply = new CountDownLatch(1); + private CountDownLatch latch = new CountDownLatch(1); - public WsClient(URI serverUri) { + WsClient(URI serverUri) { super(serverUri); } @@ -43,14 +45,19 @@ public class WsClient extends WebSocketClient { @Override public void onMessage(String message) { - try { - WsTelemetryResponse response = mapper.readValue(message, WsTelemetryResponse.class); - if (!response.getData().isEmpty()) { - this.message = response; - latch.countDown(); + if (!firstReplyReceived) { + firstReplyReceived = true; + firstReply.countDown(); + } else { + try { + WsTelemetryResponse response = mapper.readValue(message, WsTelemetryResponse.class); + if (!response.getData().isEmpty()) { + this.message = response; + latch.countDown(); + } + } catch (IOException e) { + log.error("ws message can't be read"); } - } catch (IOException e) { - log.error("ws message can't be read"); } } @@ -73,4 +80,13 @@ public class WsClient extends WebSocketClient { } return null; } + + void waitForFirstReply() { + try { + firstReply.await(10, TimeUnit.SECONDS); + } catch (InterruptedException e) { + log.error("Timeout, ws message wasn't received"); + throw new RuntimeException(e); + } + } } diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java index d889d2c52e..9918963a1c 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/connectivity/MqttClientTest.java @@ -28,6 +28,9 @@ import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.RandomStringUtils; import org.junit.*; +import org.junit.rules.TestRule; +import org.junit.rules.TestWatcher; +import org.junit.runner.Description; import org.springframework.core.ParameterizedTypeReference; import org.springframework.http.HttpMethod; import org.springframework.http.ResponseEntity; @@ -65,6 +68,7 @@ public class MqttClientTest extends AbstractContainerTest { MqttClient mqttClient = getMqttClient(deviceCredentials, null); mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload().toString().getBytes())); WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); + log.info("Received telemetry: {}", actualLatestTelemetry); wsClient.closeBlocking(); Assert.assertEquals(4, actualLatestTelemetry.getData().size()); @@ -91,6 +95,7 @@ public class MqttClientTest extends AbstractContainerTest { MqttClient mqttClient = getMqttClient(deviceCredentials, null); mqttClient.publish("v1/devices/me/telemetry", Unpooled.wrappedBuffer(createPayload(ts).toString().getBytes())); WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); + log.info("Received telemetry: {}", actualLatestTelemetry); wsClient.closeBlocking(); Assert.assertEquals(4, actualLatestTelemetry.getData().size()); @@ -120,6 +125,7 @@ public class MqttClientTest extends AbstractContainerTest { clientAttributes.addProperty("attr4", 73); mqttClient.publish("v1/devices/me/attributes", Unpooled.wrappedBuffer(clientAttributes.toString().getBytes())); WsTelemetryResponse actualLatestTelemetry = wsClient.getLastMessage(); + log.info("Received telemetry: {}", actualLatestTelemetry); wsClient.closeBlocking(); Assert.assertEquals(4, actualLatestTelemetry.getData().size()); @@ -168,6 +174,7 @@ public class MqttClientTest extends AbstractContainerTest { mqttClient.publish("v1/devices/me/attributes/request/" + new Random().nextInt(100), Unpooled.wrappedBuffer(request.toString().getBytes())); MqttEvent event = listener.getEvents().poll(10, TimeUnit.SECONDS); AttributesResponse attributes = mapper.readValue(Objects.requireNonNull(event).getMessage(), AttributesResponse.class); + log.info("Received telemetry: {}", attributes); Assert.assertEquals(1, attributes.getClient().size()); Assert.assertEquals(clientAttributeValue, attributes.getClient().get("clientAttr")); @@ -281,6 +288,7 @@ public class MqttClientTest extends AbstractContainerTest { // Create a new root rule chain RuleChainId ruleChainId = createRootRuleChainForRpcResponse(); + TimeUnit.SECONDS.sleep(3); // Send the request to the server JsonObject clientRequest = new JsonObject(); clientRequest.addProperty("method", "getResponse"); @@ -360,12 +368,12 @@ public class MqttClientTest extends AbstractContainerTest { return defaultRuleChain.get().getId(); } - private MqttClient getMqttClient(DeviceCredentials deviceCredentials, MqttMessageListener listener) throws InterruptedException { + private MqttClient getMqttClient(DeviceCredentials deviceCredentials, MqttMessageListener listener) throws InterruptedException, ExecutionException { MqttClientConfig clientConfig = new MqttClientConfig(); clientConfig.setClientId("MQTT client from test"); clientConfig.setUsername(deviceCredentials.getCredentialsId()); MqttClient mqttClient = MqttClient.create(clientConfig, listener); - mqttClient.connect("localhost", 1883).sync(); + mqttClient.connect("localhost", 1883).get(); return mqttClient; }