diff --git a/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java b/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java index 9f82b6ca2b..c3a9d9817a 100644 --- a/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/rules/lifecycle/AbstractRuleEngineLifecycleIntegrationTest.java @@ -15,30 +15,24 @@ */ package org.thingsboard.server.rules.lifecycle; -import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; +import org.awaitility.Awaitility; import org.junit.After; import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.mockito.Mockito; -import org.mockito.stubbing.Answer; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.test.util.ReflectionTestUtils; import org.thingsboard.rule.engine.metadata.TbGetAttributesNodeConfiguration; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Event; -import org.thingsboard.server.common.data.Tenant; -import org.thingsboard.server.common.data.User; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; -import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.rule.RuleNode; -import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; import org.thingsboard.server.common.msg.queue.QueueToRuleEngineMsg; @@ -46,16 +40,13 @@ import org.thingsboard.server.common.msg.queue.TbMsgCallback; import org.thingsboard.server.controller.AbstractRuleEngineControllerTest; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.event.EventService; -import org.thingsboard.server.queue.memory.InMemoryStorage; import java.util.Collections; import java.util.List; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; -import static org.awaitility.Awaitility.await; -import static org.mockito.Mockito.spy; -import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import static java.util.concurrent.TimeUnit.MILLISECONDS; /** * @author Valerii Sosliuk @@ -63,9 +54,6 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers. @Slf4j public abstract class AbstractRuleEngineLifecycleIntegrationTest extends AbstractRuleEngineControllerTest { - protected Tenant savedTenant; - protected User tenantAdmin; - @Autowired protected ActorSystemContext actorSystem; @@ -75,50 +63,14 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac @Autowired protected EventService eventService; - @Autowired - protected InMemoryStorage storage; - @Before public void beforeTest() throws Exception { - - EventService spyEventService = spy(eventService); - - Mockito.doAnswer((Answer>) invocation -> { - Object[] args = invocation.getArguments(); - Event event = (Event) args[0]; - ListenableFuture future = eventService.saveAsync(event); - try { - future.get(); - } catch (Exception e) {} - return future; - }).when(spyEventService).saveAsync(Mockito.any(Event.class)); - - ReflectionTestUtils.setField(actorSystem, "eventService", spyEventService); - - loginSysAdmin(); - - Tenant tenant = new Tenant(); - tenant.setTitle("My tenant"); - savedTenant = doPost("/api/tenant", tenant, Tenant.class); - Assert.assertNotNull(savedTenant); - ruleChainService.deleteRuleChainsByTenantId(savedTenant.getId()); - - tenantAdmin = new User(); - tenantAdmin.setAuthority(Authority.TENANT_ADMIN); - tenantAdmin.setTenantId(savedTenant.getId()); - tenantAdmin.setEmail("tenant2@thingsboard.org"); - tenantAdmin.setFirstName("Joe"); - tenantAdmin.setLastName("Downs"); - - createUserAndLogin(tenantAdmin, "testPassword1"); + loginTenantAdmin(); + ruleChainService.deleteRuleChainsByTenantId(tenantId); } @After public void afterTest() throws Exception { - loginSysAdmin(); - if (savedTenant != null) { - doDelete("/api/tenant/" + savedTenant.getId().getId().toString()).andExpect(status().isOk()); - } } @Test @@ -126,7 +78,7 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac // Creating Rule Chain RuleChain ruleChain = new RuleChain(); ruleChain.setName("Simple Rule Chain"); - ruleChain.setTenantId(savedTenant.getId()); + ruleChain.setTenantId(tenantId); ruleChain.setRoot(true); ruleChain.setDebugMode(true); ruleChain = saveRuleChain(ruleChain); @@ -149,8 +101,21 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac metaData = saveRuleChainMetaData(metaData); Assert.assertNotNull(metaData); - ruleChain = getRuleChain(ruleChain.getId()); - Assert.assertNotNull(ruleChain.getFirstRuleNodeId()); + final RuleChain ruleChainFinal = getRuleChain(ruleChain.getId()); + Assert.assertNotNull(ruleChainFinal.getFirstRuleNodeId()); + + //TODO find out why RULE_NODE update event did not appear all the time +// List rcEvents = Awaitility.await("get debug by rule chain") +// .pollInterval(10, MILLISECONDS) +// .atMost(TIMEOUT, TimeUnit.SECONDS) +// .until(() -> { +// List debugEvents = getDebugEvents(tenantId, ruleChainFinal.getFirstRuleNodeId(), 1000) +// .getData().stream().filter(x->true).collect(Collectors.toList()); +// log.warn("filtered debug events [{}]", debugEvents.size()); +// debugEvents.forEach((e) -> log.warn("event: {}", e)); +// return debugEvents; +// }, +// x -> x.size() >= 2); // Saving the device Device device = new Device(); @@ -158,35 +123,45 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac device.setType("default"); device = doPost("/api/device", device, Device.class); + log.warn("before update attr"); attributesService.save(device.getTenantId(), device.getId(), DataConstants.SERVER_SCOPE, - Collections.singletonList(new BaseAttributeKvEntry(new StringDataEntry("serverAttributeKey", "serverAttributeValue"), System.currentTimeMillis()))); - - await("total inMemory queue lag is empty").atMost(30, TimeUnit.SECONDS) - .until(() -> storage.getLagTotal() == 0); - Thread.sleep(1000); - + Collections.singletonList(new BaseAttributeKvEntry(new StringDataEntry("serverAttributeKey", "serverAttributeValue"), System.currentTimeMillis()))) + .get(TIMEOUT, TimeUnit.SECONDS); + log.warn("attr updated"); TbMsgCallback tbMsgCallback = Mockito.mock(TbMsgCallback.class); Mockito.when(tbMsgCallback.isMsgValid()).thenReturn(true); TbMsg tbMsg = TbMsg.newMsg("CUSTOM", device.getId(), new TbMsgMetaData(), "{}", tbMsgCallback); - QueueToRuleEngineMsg qMsg = new QueueToRuleEngineMsg(savedTenant.getId(), tbMsg, null, null); + QueueToRuleEngineMsg qMsg = new QueueToRuleEngineMsg(tenantId, tbMsg, null, null); // Pushing Message to the system actorSystem.tell(qMsg); - Mockito.verify(tbMsgCallback, Mockito.timeout(10000)).onSuccess(); - - - PageData eventsPage = getDebugEvents(savedTenant.getId(), ruleChain.getFirstRuleNodeId(), 1000); - List events = eventsPage.getData().stream().filter(filterByCustomEvent()).collect(Collectors.toList()); - - Assert.assertEquals(2, events.size()); + log.warn("awaiting tbMsgCallback"); + Mockito.verify(tbMsgCallback, Mockito.timeout(TimeUnit.SECONDS.toMillis(TIMEOUT))).onSuccess(); + log.warn("awaiting events"); + List events = Awaitility.await("get debug by custom event") + .pollInterval(10, MILLISECONDS) + .atMost(TIMEOUT, TimeUnit.SECONDS) + .until(() -> { + List debugEvents = getDebugEvents(tenantId, ruleChainFinal.getFirstRuleNodeId(), 1000) + .getData().stream().filter(filterByCustomEvent()).collect(Collectors.toList()); + log.warn("filtered debug events [{}]", debugEvents.size()); + debugEvents.forEach((e) -> log.warn("event: {}", e)); + return debugEvents; + }, + x -> x.size() == 2); + log.warn("asserting.."); Event inEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); - Assert.assertEquals(ruleChain.getFirstRuleNodeId(), inEvent.getEntityId()); + Assert.assertEquals(ruleChainFinal.getFirstRuleNodeId(), inEvent.getEntityId()); Assert.assertEquals(device.getId().getId().toString(), inEvent.getBody().get("entityId").asText()); Event outEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.OUT)).findFirst().get(); - Assert.assertEquals(ruleChain.getFirstRuleNodeId(), outEvent.getEntityId()); + Assert.assertEquals(ruleChainFinal.getFirstRuleNodeId(), outEvent.getEntityId()); Assert.assertEquals(device.getId().getId().toString(), outEvent.getBody().get("entityId").asText()); + log.warn("OUT event {}", outEvent); + log.warn("OUT event metadata {}", getMetadata(outEvent)); + + Assert.assertNotNull("metadata has ss_serverAttributeKey", getMetadata(outEvent).get("ss_serverAttributeKey")); Assert.assertEquals("serverAttributeValue", getMetadata(outEvent).get("ss_serverAttributeKey").asText()); }