Browse Source

code coverage for processors and telemetry constructor

pull/2436/head
Bohdan Smetaniuk 6 years ago
parent
commit
9e21525158
  1. 284
      application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java
  2. 25
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java
  3. 29
      application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java
  4. 2
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java

284
application/src/test/java/org/thingsboard/server/edge/BaseEdgeTest.java

@ -15,7 +15,11 @@
*/
package org.thingsboard.server.edge;
import com.datastax.driver.core.utils.UUIDs;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.google.gson.JsonObject;
import lombok.extern.slf4j.Slf4j;
import org.junit.After;
import org.junit.Assert;
@ -24,6 +28,7 @@ import org.junit.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.thingsboard.server.common.data.Dashboard;
import org.thingsboard.server.common.data.DashboardInfo;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.Tenant;
@ -33,7 +38,11 @@ import org.thingsboard.server.common.data.alarm.AlarmInfo;
import org.thingsboard.server.common.data.alarm.AlarmSeverity;
import org.thingsboard.server.common.data.alarm.AlarmStatus;
import org.thingsboard.server.common.data.asset.Asset;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.edge.Edge;
import org.thingsboard.server.common.data.edge.EdgeEvent;
import org.thingsboard.server.common.data.edge.EdgeEventType;
import org.thingsboard.server.common.data.id.EdgeId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.TextPageLink;
import org.thingsboard.server.common.data.page.TimePageData;
@ -42,17 +51,24 @@ import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.common.data.rule.RuleChain;
import org.thingsboard.server.common.data.rule.RuleChainType;
import org.thingsboard.server.common.data.security.Authority;
import org.thingsboard.server.common.transport.adaptor.JsonConverter;
import org.thingsboard.server.controller.AbstractControllerTest;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.edge.EdgeEventService;;
import org.thingsboard.server.edge.imitator.EdgeImitator;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.DeviceUpdateMsg;
import org.thingsboard.server.gen.edge.EdgeConfiguration;
import org.thingsboard.server.gen.edge.EntityDataProto;
import org.thingsboard.server.gen.edge.RelationUpdateMsg;
import org.thingsboard.server.gen.edge.UpdateMsgType;
import org.thingsboard.server.gen.edge.UplinkMsg;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;
@ -67,6 +83,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
private EdgeImitator edgeImitator;
private Edge edge;
@Autowired
private EdgeEventService edgeEventService;
@Before
public void beforeTest() throws Exception {
loginSysAdmin();
@ -114,6 +133,9 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
testDashboards();
testRelations();
testAlarms();
testTimeseries();
testAttributes();
testSendMessagesToCloud();
}
private void testReceivedInitialData() throws Exception {
@ -376,6 +398,251 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
log.info("Alarms tested successfully");
}
private void testTimeseries() throws Exception {
log.info("Testing timeseries");
ObjectMapper mapper = new ObjectMapper();
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Assert.assertEquals(1, edgeDevices.size());
Device device = edgeDevices.get(0);
Assert.assertEquals("Edge Device 1", device.getName());
String timeseriesData = "{\"data\":{\"temperature\":25},\"ts\":" + System.currentTimeMillis() + "}";
JsonNode timeseriesEntityData = mapper.readTree(timeseriesData);
EdgeEvent edgeEvent1 = constructEdgeEvent(tenantId, edge.getId(), ActionType.TIMESERIES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, timeseriesEntityData);
edgeImitator.getStorage().expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent1);
edgeImitator.getStorage().waitForMessages();
EntityDataProto latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg();
Assert.assertNotNull(latestEntityDataMsg);
UUID uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(device.getId().getId(), uuid);
Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertTrue(latestEntityDataMsg.hasPostTelemetryMsg());
TransportProtos.PostTelemetryMsg postTelemetryMsg = latestEntityDataMsg.getPostTelemetryMsg();
Assert.assertEquals(1, postTelemetryMsg.getTsKvListCount());
TransportProtos.TsKvListProto tsKvListProto = postTelemetryMsg.getTsKvList(0);
Assert.assertEquals(timeseriesEntityData.get("ts").asLong(), tsKvListProto.getTs());
Assert.assertEquals(1, tsKvListProto.getKvCount());
TransportProtos.KeyValueProto keyValueProto = tsKvListProto.getKv(0);
Assert.assertEquals("temperature", keyValueProto.getKey());
Assert.assertEquals(25, keyValueProto.getLongV());
edgeImitator.getStorage().setLatestEntityDataMsg(null);
log.info("Timeseries tested successfully");
}
private void testAttributes() throws Exception {
log.info("Testing attributes");
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Assert.assertEquals(1, edgeDevices.size());
Device device = edgeDevices.get(0);
Assert.assertEquals("Edge Device 1", device.getName());
String attributesData = "{\"scope\":\"SERVER_SCOPE\",\"kv\":{\"test\":\"test\"}}";
JsonNode attributesEntityData = mapper.readTree(attributesData);
EdgeEvent edgeEvent2 = constructEdgeEvent(tenantId, edge.getId(), ActionType.ATTRIBUTES_UPDATED, device.getId().getId(), EdgeEventType.DEVICE, attributesEntityData);
edgeImitator.getStorage().expectMessageAmount(1);
edgeEventService.saveAsync(edgeEvent2);
edgeImitator.getStorage().waitForMessages();
EntityDataProto latestEntityDataMsg = edgeImitator.getStorage().getLatestEntityDataMsg();
Assert.assertNotNull(latestEntityDataMsg);
UUID uuid = new UUID(latestEntityDataMsg.getEntityIdMSB(), latestEntityDataMsg.getEntityIdLSB());
Assert.assertEquals(device.getId().getId(), uuid);
Assert.assertEquals(device.getId().getEntityType().name(), latestEntityDataMsg.getEntityType());
Assert.assertEquals(attributesEntityData.get("scope").asText(), latestEntityDataMsg.getPostAttributeScope());
Assert.assertTrue(latestEntityDataMsg.hasPostAttributesMsg());
TransportProtos.PostAttributeMsg postAttributeMsg = latestEntityDataMsg.getPostAttributesMsg();
Assert.assertEquals(1, postAttributeMsg.getKvCount());
TransportProtos.KeyValueProto keyValueProto = postAttributeMsg.getKv(0);
Assert.assertEquals("test", keyValueProto.getKey());
Assert.assertEquals("test", keyValueProto.getStringV());
edgeImitator.getStorage().setLatestEntityDataMsg(null);
log.info("Attributes tested successfully");
}
private void testSendMessagesToCloud() throws Exception {
log.info("Sending messages to cloud");
sendDevice();
sendAlarm();
sendTelemetry();
sendRelation();
sendDeleteDeviceOnEdge();
log.info("Messages were sent successfully");
}
private void sendDevice() throws Exception {
UUID uuid = UUIDs.timeBased();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
DeviceUpdateMsg.Builder deviceUpdateMsgBuilder = DeviceUpdateMsg.newBuilder();
deviceUpdateMsgBuilder.setIdMSB(uuid.getMostSignificantBits());
deviceUpdateMsgBuilder.setIdLSB(uuid.getLeastSignificantBits());
deviceUpdateMsgBuilder.setName("Edge Device 2");
deviceUpdateMsgBuilder.setType("test");
deviceUpdateMsgBuilder.setMsgType(UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE);
builder.addDeviceUpdateMsg(deviceUpdateMsgBuilder.build());
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build());
edgeImitator.waitForResponses();
Device device = doGet("/api/device/" + uuid.toString(), Device.class);
Assert.assertNotNull(device);
Assert.assertEquals("Edge Device 2", device.getName());
}
private void sendAlarm() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
AlarmUpdateMsg.Builder alarmUpdateMgBuilder = AlarmUpdateMsg.newBuilder();
alarmUpdateMgBuilder.setName("alarm from edge");
alarmUpdateMgBuilder.setStatus(AlarmStatus.ACTIVE_UNACK.name());
alarmUpdateMgBuilder.setSeverity(AlarmSeverity.CRITICAL.name());
alarmUpdateMgBuilder.setOriginatorName(device.getName());
alarmUpdateMgBuilder.setOriginatorType(EntityType.DEVICE.name());
builder.addAlarmUpdateMsg(alarmUpdateMgBuilder.build());
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build());
edgeImitator.waitForResponses();
List<AlarmInfo> alarms = doGetTypedWithPageLink("/api/alarm/{entityType}/{entityId}?",
new TypeReference<TimePageData<AlarmInfo>>() {},
new TextPageLink(100), device.getId().getEntityType().name(), device.getId().getId().toString())
.getData();
for (AlarmInfo alarmInfo: alarms) {
log.info(String.valueOf(alarmInfo));
}
Optional<AlarmInfo> foundAlarm = alarms.stream().filter(alarm -> alarm.getType().equals("alarm from edge")).findAny();
Assert.assertTrue(foundAlarm.isPresent());
AlarmInfo alarmInfo = foundAlarm.get();
Assert.assertEquals(device.getId(), alarmInfo.getOriginator());
Assert.assertEquals(AlarmStatus.ACTIVE_UNACK, alarmInfo.getStatus());
Assert.assertEquals(AlarmSeverity.CRITICAL, alarmInfo.getSeverity());
}
private void sendRelation() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice1 = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 1")).findAny();
Assert.assertTrue(foundDevice1.isPresent());
Device device1 = foundDevice1.get();
Optional<Device> foundDevice2 = edgeDevices.stream().filter(device2 -> device2.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice2.isPresent());
Device device2 = foundDevice2.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
RelationUpdateMsg.Builder relationUpdateMsgBuilder = RelationUpdateMsg.newBuilder();
relationUpdateMsgBuilder.setType("test");
relationUpdateMsgBuilder.setTypeGroup(RelationTypeGroup.COMMON.name());
relationUpdateMsgBuilder.setToIdMSB(device1.getId().getId().getMostSignificantBits());
relationUpdateMsgBuilder.setToIdLSB(device1.getId().getId().getLeastSignificantBits());
relationUpdateMsgBuilder.setToEntityType(device1.getId().getEntityType().name());
relationUpdateMsgBuilder.setFromIdMSB(device2.getId().getId().getMostSignificantBits());
relationUpdateMsgBuilder.setFromIdLSB(device2.getId().getId().getLeastSignificantBits());
relationUpdateMsgBuilder.setFromEntityType(device2.getId().getEntityType().name());
relationUpdateMsgBuilder.setAdditionalInfo("{}");
builder.addRelationUpdateMsg(relationUpdateMsgBuilder.build());
UplinkMsg msg = builder.build();
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(msg);
edgeImitator.waitForResponses();
EntityRelation relation = doGet("/api/relation?" +
"&fromId=" + device2.getId().getId().toString() +
"&fromType=" + device2.getId().getEntityType().name() +
"&relationType=" + "test" +
"&relationTypeGroup=" + RelationTypeGroup.COMMON.name() +
"&toId=" + device1.getId().getId().toString() +
"&toType=" + device1.getId().getEntityType().name(), EntityRelation.class);
Assert.assertNotNull(relation);
}
private void sendTelemetry() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get();
edgeImitator.expectResponsesAmount(2);
JsonObject data = new JsonObject();
String timeseriesKey = "key";
String timeseriesValue = "25";
data.addProperty(timeseriesKey, timeseriesValue);
UplinkMsg.Builder builder1 = UplinkMsg.newBuilder();
EntityDataProto.Builder entityDataBuilder = EntityDataProto.newBuilder();
entityDataBuilder.setPostTelemetryMsg(JsonConverter.convertToTelemetryProto(data, System.currentTimeMillis()));
entityDataBuilder.setEntityType(device.getId().getEntityType().name());
entityDataBuilder.setEntityIdMSB(device.getUuidId().getMostSignificantBits());
entityDataBuilder.setEntityIdLSB(device.getUuidId().getLeastSignificantBits());
builder1.addEntityData(entityDataBuilder.build());
edgeImitator.sendUplinkMsg(builder1.build());
JsonObject attributesData = new JsonObject();
String attributesKey = "test_attr";
String attributesValue = "test_value";
attributesData.addProperty(attributesKey, attributesValue);
UplinkMsg.Builder builder2 = UplinkMsg.newBuilder();
EntityDataProto.Builder entityDataBuilder2 = EntityDataProto.newBuilder();
entityDataBuilder2.setEntityType(device.getId().getEntityType().name());
entityDataBuilder2.setEntityIdMSB(device.getId().getId().getMostSignificantBits());
entityDataBuilder2.setEntityIdLSB(device.getId().getId().getLeastSignificantBits());
entityDataBuilder2.setPostAttributesMsg(JsonConverter.convertToAttributesProto(attributesData));
entityDataBuilder2.setPostAttributeScope(DataConstants.SERVER_SCOPE);
builder2.addEntityData(entityDataBuilder2.build());
edgeImitator.sendUplinkMsg(builder2.build());
edgeImitator.waitForResponses();
Thread.sleep(1000);
Map<String, List<Map<String, String>>> timeseries = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getUuidId() + "/values/timeseries?keys=" + timeseriesKey, Map.class);
Assert.assertTrue(timeseries.containsKey(timeseriesKey));
Assert.assertEquals(1, timeseries.get(timeseriesKey).size());
Assert.assertEquals(timeseriesValue, timeseries.get(timeseriesKey).get(0).get("value"));
List<Map<String, String>> attributes = doGetAsync("/api/plugins/telemetry/DEVICE/" + device.getId() + "/values/attributes/" + DataConstants.SERVER_SCOPE, List.class);
Assert.assertEquals(1, attributes.size());
Assert.assertEquals(attributes.get(0).get("key"), attributesKey);
Assert.assertEquals(attributes.get(0).get("value"), attributesValue);
}
private void sendDeleteDeviceOnEdge() throws Exception {
List<Device> edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {}, new TextPageLink(100)).getData();
Optional<Device> foundDevice = edgeDevices.stream().filter(device1 -> device1.getName().equals("Edge Device 2")).findAny();
Assert.assertTrue(foundDevice.isPresent());
Device device = foundDevice.get();
UplinkMsg.Builder builder = UplinkMsg.newBuilder();
DeviceUpdateMsg.Builder deviceDeleteMsgBuilder = DeviceUpdateMsg.newBuilder();
deviceDeleteMsgBuilder.setMsgType(UpdateMsgType.ENTITY_DELETED_RPC_MESSAGE);
deviceDeleteMsgBuilder.setIdMSB(device.getId().getId().getMostSignificantBits());
deviceDeleteMsgBuilder.setIdLSB(device.getId().getId().getLeastSignificantBits());
builder.addDeviceUpdateMsg(deviceDeleteMsgBuilder.build());
edgeImitator.expectResponsesAmount(1);
edgeImitator.sendUplinkMsg(builder.build());
edgeImitator.waitForResponses();
device = doGet("/api/device/" + device.getId().getId().toString(), Device.class);
Assert.assertNotNull(device);
edgeDevices = doGetTypedWithPageLink("/api/edge/" + edge.getId().getId().toString() + "/devices?",
new TypeReference<TimePageData<Device>>() {
}, new TextPageLink(100)).getData();
Assert.assertFalse(edgeDevices.contains(device));
}
private void installation() throws Exception {
edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class);
@ -413,4 +680,15 @@ abstract public class BaseEdgeTest extends AbstractControllerTest {
doDelete("/api/edge/" + edge.getId().getId().toString())
.andExpect(status().isOk());
}
private EdgeEvent constructEdgeEvent(TenantId tenantId, EdgeId edgeId, ActionType edgeEventAction, UUID entityId, EdgeEventType edgeEventType, JsonNode entityBody) {
EdgeEvent edgeEvent = new EdgeEvent();
edgeEvent.setEdgeId(edgeId);
edgeEvent.setTenantId(tenantId);
edgeEvent.setAction(edgeEventAction.name());
edgeEvent.setEntityId(entityId);
edgeEvent.setType(edgeEventType);
edgeEvent.setBody(entityBody);
return edgeEvent;
}
}

25
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java

@ -32,14 +32,18 @@ import org.thingsboard.server.gen.edge.DeviceUpdateMsg;
import org.thingsboard.server.gen.edge.DownlinkMsg;
import org.thingsboard.server.gen.edge.DownlinkResponseMsg;
import org.thingsboard.server.gen.edge.EdgeConfiguration;
import org.thingsboard.server.gen.edge.EntityDataProto;
import org.thingsboard.server.gen.edge.RelationUpdateMsg;
import org.thingsboard.server.gen.edge.RuleChainUpdateMsg;
import org.thingsboard.server.gen.edge.UplinkMsg;
import org.thingsboard.server.gen.edge.UplinkResponseMsg;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
@Slf4j
public class EdgeImitator {
@ -49,6 +53,8 @@ public class EdgeImitator {
private EdgeRpcClient edgeRpcClient;
private CountDownLatch responsesLatch;
@Getter
private EdgeStorage storage;
@ -56,10 +62,12 @@ public class EdgeImitator {
public EdgeImitator(String host, int port, String routingKey, String routingSecret) throws NoSuchFieldException, IllegalAccessException {
edgeRpcClient = new EdgeGrpcClient();
storage = new EdgeStorage();
responsesLatch = new CountDownLatch(0);
this.routingKey = routingKey;
this.routingSecret = routingSecret;
setEdgeCredentials("rpcHost", host);
setEdgeCredentials("rpcPort", port);
setEdgeCredentials("keepAliveTimeSec", 300);
}
private void setEdgeCredentials(String fieldName, Object value) throws NoSuchFieldException, IllegalAccessException {
@ -83,8 +91,13 @@ public class EdgeImitator {
edgeRpcClient.disconnect(false);
}
public void sendUplinkMsg(UplinkMsg uplinkMsg) {
edgeRpcClient.sendUplinkMsg(uplinkMsg);
}
private void onUplinkResponse(UplinkResponseMsg msg) {
log.info("onUplinkResponse: {}", msg);
responsesLatch.countDown();
}
private void onEdgeUpdate(EdgeConfiguration edgeConfiguration) {
@ -144,7 +157,19 @@ public class EdgeImitator {
result.add(storage.processAlarm(alarmUpdateMsg));
}
}
if (downlinkMsg.getEntityDataList() != null && !downlinkMsg.getEntityDataList().isEmpty()) {
for (EntityDataProto entityData: downlinkMsg.getEntityDataList()) {
result.add(storage.processEntityData(entityData));
}
}
return Futures.allAsList(result);
}
public void waitForResponses() throws InterruptedException { responsesLatch.await(5, TimeUnit.SECONDS);
}
public void expectResponsesAmount(int messageAmount) {
responsesLatch = new CountDownLatch(messageAmount);
}
}

29
application/src/test/java/org/thingsboard/server/edge/imitator/EdgeStorage.java

@ -27,6 +27,7 @@ import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.RelationTypeGroup;
import org.thingsboard.server.gen.edge.AlarmUpdateMsg;
import org.thingsboard.server.gen.edge.EdgeConfiguration;
import org.thingsboard.server.gen.edge.EntityDataProto;
import org.thingsboard.server.gen.edge.RelationUpdateMsg;
import org.thingsboard.server.gen.edge.UpdateMsgType;
@ -47,17 +48,20 @@ public class EdgeStorage {
private EdgeConfiguration configuration;
private CountDownLatch latch;
private CountDownLatch messagesLatch;
private Map<UUID, EntityType> entities;
private Map<String, AlarmStatus> alarms;
private List<EntityRelation> relations;
private EntityDataProto latestEntityDataMsg;
public EdgeStorage() {
latch = new CountDownLatch(0);
messagesLatch = new CountDownLatch(0);
entities = new HashMap<>();
alarms = new HashMap<>();
relations = new ArrayList<>();
latestEntityDataMsg = null;
}
public ListenableFuture<Void> processEntity(UpdateMsgType msgType, EntityType type, UUID uuid) {
@ -65,11 +69,11 @@ public class EdgeStorage {
case ENTITY_CREATED_RPC_MESSAGE:
case ENTITY_UPDATED_RPC_MESSAGE:
entities.put(uuid, type);
latch.countDown();
messagesLatch.countDown();
break;
case ENTITY_DELETED_RPC_MESSAGE:
if (entities.remove(uuid) != null) {
latch.countDown();
messagesLatch.countDown();
}
break;
}
@ -93,7 +97,7 @@ public class EdgeStorage {
break;
}
if (result) {
latch.countDown();
messagesLatch.countDown();
}
return Futures.immediateFuture(null);
}
@ -105,17 +109,23 @@ public class EdgeStorage {
case ALARM_ACK_RPC_MESSAGE:
case ALARM_CLEAR_RPC_MESSAGE:
alarms.put(alarmMsg.getType(), AlarmStatus.valueOf(alarmMsg.getStatus()));
latch.countDown();
messagesLatch.countDown();
break;
case ENTITY_DELETED_RPC_MESSAGE:
if (alarms.remove(alarmMsg.getName()) != null) {
latch.countDown();
messagesLatch.countDown();
}
break;
}
return Futures.immediateFuture(null);
}
public ListenableFuture<Void> processEntityData(EntityDataProto entityData) {
latestEntityDataMsg = entityData;
messagesLatch.countDown();
return Futures.immediateFuture(null);
}
public Set<UUID> getEntitiesByType(EntityType type) {
return entities.entrySet().stream()
.filter(entry -> entry.getValue().equals(type))
@ -123,11 +133,10 @@ public class EdgeStorage {
}
public void waitForMessages() throws InterruptedException {
latch.await(5, TimeUnit.SECONDS);
messagesLatch.await(5, TimeUnit.SECONDS);
}
public void expectMessageAmount(int messageAmount) {
latch = new CountDownLatch(messageAmount);
messagesLatch = new CountDownLatch(messageAmount);
}
}

2
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java

@ -66,7 +66,7 @@ public class TbMsgAttributesNode implements TbNode {
if (StringUtils.isEmpty(msg.getMetaData().getValue(SCOPE))) {
msg.getMetaData().putValue(SCOPE, config.getScope());
}
ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), msg.getOriginator(), msg.getMetaData().getValue(SCOPE),
ctx.getTelemetryService().saveAndNotify(ctx.getTenantId(), msg.getOriginator(), config.getScope(),
new ArrayList<>(attributes), new TelemetryNodeCallback(ctx, msg));
}

Loading…
Cancel
Save