From a720009b55678d30ab5f685ba3acc65934011b2d Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 22 Jul 2022 18:53:41 +0300 Subject: [PATCH] Initial DB structure --- .../server/actors/ActorSystemContext.java | 159 +++++----- .../server/actors/stats/StatsActor.java | 18 +- .../server/controller/EventController.java | 8 +- .../controller/RuleChainController.java | 7 +- .../server/actors/stats/StatsActorTest.java | 3 +- .../AbstractRuleEngineControllerTest.java | 12 +- ...AbstractRuleEngineFlowIntegrationTest.java | 19 +- ...actRuleEngineLifecycleIntegrationTest.java | 18 +- .../server/dao/event/EventService.java | 13 +- .../data/{Event.java => EventInfo.java} | 8 +- .../server/common/data/event/ErrorEvent.java | 48 +++ .../server/common/data/event/Event.java | 42 +++ .../server/common/data/event/EventType.java | 12 +- .../common/data/event/LifecycleEvent.java | 51 ++++ .../data/event/RuleChainDebugEvent.java | 50 +++ .../common/data/event/RuleNodeDebugEvent.java | 75 +++++ .../common/data/event/StatisticsEvent.java | 47 +++ .../server/dao/event/BaseEventService.java | 59 +++- .../server/dao/event/EventDao.java | 13 +- .../server/dao/model/sql/EventEntity.java | 9 +- .../service/validator/EventDataValidator.java | 12 +- .../dao/sql/event/EventInsertRepository.java | 201 ++++++++++-- .../server/dao/sql/event/JpaBaseEventDao.java | 287 ++++++++++-------- .../insert/sql/SqlPartitioningRepository.java | 4 +- .../dao/sqlts/sql/JpaSqlTimeseriesDao.java | 2 +- .../server/dao/timeseries/SqlPartition.java | 10 +- .../main/resources/sql/schema-entities.sql | 68 ++++- .../dao/service/AbstractServiceTest.java | 20 +- .../service/event/BaseEventServiceTest.java | 50 +-- .../dao/sql/event/JpaBaseEventDaoTest.java | 249 ++++++++------- .../thingsboard/rest/client/RestClient.java | 10 +- 31 files changed, 1074 insertions(+), 510 deletions(-) rename common/data/src/main/java/org/thingsboard/server/common/data/{Event.java => EventInfo.java} (93%) create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/event/ErrorEvent.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/event/Event.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/event/LifecycleEvent.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/event/RuleChainDebugEvent.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/event/RuleNodeDebugEvent.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/event/StatisticsEvent.java diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java index a38de7b839..e6f207d960 100644 --- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java +++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java @@ -37,10 +37,11 @@ import org.thingsboard.rule.engine.api.sms.SmsSenderFactory; import org.thingsboard.server.actors.service.ActorService; import org.thingsboard.server.actors.tenant.DebugTbRateLimits; import org.thingsboard.server.cluster.TbClusterService; -import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.event.ErrorEvent; +import org.thingsboard.server.common.data.event.LifecycleEvent; +import org.thingsboard.server.common.data.event.RuleChainDebugEvent; +import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.QueueId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.plugin.ComponentLifecycleEvent; import org.thingsboard.server.common.msg.TbActorMsg; @@ -112,6 +113,29 @@ import java.util.concurrent.TimeUnit; @Component public class ActorSystemContext { + private static final FutureCallback RULE_CHAIN_DEBUG_EVENT_ERROR_CALLBACK = new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Void event) { + + } + + @Override + public void onFailure(Throwable th) { + log.error("Could not save debug Event for Rule Chain", th); + } + }; + private static final FutureCallback RULE_NODE_DEBUG_EVENT_ERROR_CALLBACK = new FutureCallback<>() { + @Override + public void onSuccess(@Nullable Void event) { + + } + + @Override + public void onFailure(Throwable th) { + log.error("Could not save debug Event for Node", th); + } + }; + protected final ObjectMapper mapper = new ObjectMapper(); private final ConcurrentMap debugPerTenantLimits = new ConcurrentHashMap<>(); @@ -462,25 +486,28 @@ public class ActorSystemContext { } public void persistError(TenantId tenantId, EntityId entityId, String method, Exception e) { - Event event = new Event(); - event.setTenantId(tenantId); - event.setEntityId(entityId); - event.setType(DataConstants.ERROR); - event.setBody(toBodyJson(serviceInfoProvider.getServiceInfo().getServiceId(), method, toString(e))); - persistEvent(event); + eventService.saveAsync(ErrorEvent.builder() + .tenantId(tenantId) + .entityId(entityId) + .serviceId(getServiceId()) + .method(method) + .error(toString(e)).build()); } public void persistLifecycleEvent(TenantId tenantId, EntityId entityId, ComponentLifecycleEvent lcEvent, Exception e) { - Event event = new Event(); - event.setTenantId(tenantId); - event.setEntityId(entityId); - event.setType(DataConstants.LC_EVENT); - event.setBody(toBodyJson(serviceInfoProvider.getServiceInfo().getServiceId(), lcEvent, Optional.ofNullable(e))); - persistEvent(event); - } + LifecycleEvent.LifecycleEventBuilder event = LifecycleEvent.builder() + .tenantId(tenantId) + .entityId(entityId) + .serviceId(getServiceId()) + .lcEventType(lcEvent.name()); + + if (e != null) { + event.success(false).error(toString(e)); + } else { + event.success(false); + } - private void persistEvent(Event event) { - eventService.saveAsync(event); + eventService.saveAsync(event.build()); } private String toString(Throwable e) { @@ -489,21 +516,6 @@ public class ActorSystemContext { return sw.toString(); } - private JsonNode toBodyJson(String serviceId, ComponentLifecycleEvent event, Optional e) { - ObjectNode node = mapper.createObjectNode().put("server", serviceId).put("event", event.name()); - if (e.isPresent()) { - node = node.put("success", false); - node = node.put("error", toString(e.get())); - } else { - node = node.put("success", true); - } - return node; - } - - private JsonNode toBodyJson(String serviceId, String method, String body) { - return mapper.createObjectNode().put("server", serviceId).put("method", method).put("error", body); - } - public TopicPartitionInfo resolve(ServiceType serviceType, TenantId tenantId, EntityId entityId) { return partitionService.resolve(serviceType, tenantId, entityId); } @@ -539,44 +551,27 @@ public class ActorSystemContext { private void persistDebugAsync(TenantId tenantId, EntityId entityId, String type, TbMsg tbMsg, String relationType, Throwable error, String failureMessage) { if (checkLimits(tenantId, tbMsg, error)) { try { - Event event = new Event(); - event.setTenantId(tenantId); - event.setEntityId(entityId); - event.setType(DataConstants.DEBUG_RULE_NODE); - - String metadata = mapper.writeValueAsString(tbMsg.getMetaData().getData()); - - ObjectNode node = mapper.createObjectNode() - .put("type", type) - .put("server", getServiceId()) - .put("entityId", tbMsg.getOriginator().getId().toString()) - .put("entityName", tbMsg.getOriginator().getEntityType().name()) - .put("msgId", tbMsg.getId().toString()) - .put("msgType", tbMsg.getType()) - .put("dataType", tbMsg.getDataType().name()) - .put("relationType", relationType) - .put("data", tbMsg.getData()) - .put("metadata", metadata); + RuleNodeDebugEvent.RuleNodeDebugEventBuilder event = RuleNodeDebugEvent.builder() + .tenantId(tenantId) + .entityId(entityId) + .serviceId(getServiceId()) + .eventType(type) + .eventEntity(tbMsg.getOriginator()) + .msgId(tbMsg.getId()) + .msgType(tbMsg.getType()) + .dataType(tbMsg.getDataType().name()) + .relationType(relationType) + .data(tbMsg.getData()) + .metadata(mapper.writeValueAsString(tbMsg.getMetaData().getData())); if (error != null) { - node = node.put("error", toString(error)); + event.error(toString(error)); } else if (failureMessage != null) { - node = node.put("error", failureMessage); + event.error(failureMessage); } - event.setBody(node); - ListenableFuture future = eventService.saveAsync(event); - Futures.addCallback(future, new FutureCallback() { - @Override - public void onSuccess(@Nullable Void event) { - - } - - @Override - public void onFailure(Throwable th) { - log.error("Could not save debug Event for Node", th); - } - }, MoreExecutors.directExecutor()); + ListenableFuture future = eventService.saveAsync(event.build()); + Futures.addCallback(future, RULE_NODE_DEBUG_EVENT_ERROR_CALLBACK, MoreExecutors.directExecutor()); } catch (IOException ex) { log.warn("Failed to persist rule node debug message", ex); } @@ -603,33 +598,17 @@ public class ActorSystemContext { } private void persistRuleChainDebugModeEvent(TenantId tenantId, EntityId entityId, Throwable error) { - Event event = new Event(); - event.setTenantId(tenantId); - event.setEntityId(entityId); - event.setType(DataConstants.DEBUG_RULE_CHAIN); - - ObjectNode node = mapper.createObjectNode() - //todo: what fields are needed here? - .put("server", getServiceId()) - .put("message", "Reached debug mode rate limit!"); - + RuleChainDebugEvent.RuleChainDebugEventBuilder event = RuleChainDebugEvent.builder() + .tenantId(tenantId) + .entityId(entityId) + .serviceId(getServiceId()) + .message("Reached debug mode rate limit!"); if (error != null) { - node = node.put("error", toString(error)); + event.error(toString(error)); } - event.setBody(node); - ListenableFuture future = eventService.saveAsync(event); - Futures.addCallback(future, new FutureCallback() { - @Override - public void onSuccess(@Nullable Void event) { - - } - - @Override - public void onFailure(Throwable th) { - log.error("Could not save debug Event for Rule Chain", th); - } - }, MoreExecutors.directExecutor()); + ListenableFuture future = eventService.saveAsync(event.build()); + Futures.addCallback(future, RULE_CHAIN_DEBUG_EVENT_ERROR_CALLBACK, MoreExecutors.directExecutor()); } public static Exception toException(Throwable error) { diff --git a/application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java b/application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java index 8cdf5cec02..fa5dedf66f 100644 --- a/application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java +++ b/application/src/main/java/org/thingsboard/server/actors/stats/StatsActor.java @@ -20,13 +20,13 @@ import com.fasterxml.jackson.databind.ObjectMapper; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.TbActor; -import org.thingsboard.server.actors.TbActorCtx; import org.thingsboard.server.actors.TbActorId; import org.thingsboard.server.actors.TbStringActorId; import org.thingsboard.server.actors.service.ContextAwareActor; import org.thingsboard.server.actors.service.ContextBasedCreator; import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; +import org.thingsboard.server.common.data.event.StatisticsEvent; import org.thingsboard.server.common.msg.MsgType; import org.thingsboard.server.common.msg.TbActorMsg; @@ -54,12 +54,14 @@ public class StatsActor extends ContextAwareActor { if (msg.isEmpty()) { return; } - Event event = new Event(); - event.setEntityId(msg.getEntityId()); - event.setTenantId(msg.getTenantId()); - event.setType(DataConstants.STATS); - event.setBody(toBodyJson(systemContext.getServiceInfoProvider().getServiceId(), msg.getMessagesProcessed(), msg.getErrorsOccurred())); - systemContext.getEventService().saveAsync(event); + systemContext.getEventService().saveAsync(StatisticsEvent.builder() + .tenantId(msg.getTenantId()) + .entityId(msg.getEntityId()) + .serviceId(systemContext.getServiceInfoProvider().getServiceId()) + .messagesProcessed(msg.getMessagesProcessed()) + .errorsOccurred(msg.getErrorsOccurred()) + .build() + ); } private JsonNode toBodyJson(String serviceId, long messagesProcessed, long errorsOccurred) { diff --git a/application/src/main/java/org/thingsboard/server/controller/EventController.java b/application/src/main/java/org/thingsboard/server/controller/EventController.java index ac56978536..dd6c05749f 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EventController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EventController.java @@ -29,7 +29,7 @@ import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.ResponseBody; import org.springframework.web.bind.annotation.ResponseStatus; import org.springframework.web.bind.annotation.RestController; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.event.EventFilter; import org.thingsboard.server.common.data.exception.ThingsboardException; import org.thingsboard.server.common.data.id.EntityId; @@ -110,7 +110,7 @@ public class EventController extends BaseController { @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/events/{entityType}/{entityId}/{eventType}", method = RequestMethod.GET) @ResponseBody - public PageData getEvents( + public PageData getEvents( @ApiParam(value = ENTITY_TYPE_PARAM_DESCRIPTION, required = true) @PathVariable(ENTITY_TYPE) String strEntityType, @ApiParam(value = ENTITY_ID_PARAM_DESCRIPTION, required = true) @@ -153,7 +153,7 @@ public class EventController extends BaseController { @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/events/{entityType}/{entityId}", method = RequestMethod.GET) @ResponseBody - public PageData getEvents( + public PageData getEvents( @ApiParam(value = ENTITY_TYPE_PARAM_DESCRIPTION, required = true) @PathVariable(ENTITY_TYPE) String strEntityType, @ApiParam(value = ENTITY_ID_PARAM_DESCRIPTION, required = true) @@ -198,7 +198,7 @@ public class EventController extends BaseController { @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") @RequestMapping(value = "/events/{entityType}/{entityId}", method = RequestMethod.POST) @ResponseBody - public PageData getEvents( + public PageData getEvents( @ApiParam(value = ENTITY_TYPE_PARAM_DESCRIPTION, required = true) @PathVariable(ENTITY_TYPE) String strEntityType, @ApiParam(value = ENTITY_ID_PARAM_DESCRIPTION, required = true) diff --git a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java index faa1711fd9..143844504a 100644 --- a/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java +++ b/application/src/main/java/org/thingsboard/server/controller/RuleChainController.java @@ -40,7 +40,7 @@ import org.thingsboard.rule.engine.api.ScriptEngine; import org.thingsboard.server.actors.ActorSystemContext; import org.thingsboard.server.actors.tenant.DebugTbRateLimits; import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.exception.ThingsboardException; @@ -70,7 +70,6 @@ import org.thingsboard.server.service.script.RuleNodeJsScriptEngine; import org.thingsboard.server.service.security.permission.Operation; import org.thingsboard.server.service.security.permission.Resource; -import java.io.IOException; import java.util.ArrayList; import java.util.List; import java.util.Map; @@ -353,10 +352,10 @@ public class RuleChainController extends BaseController { RuleNodeId ruleNodeId = new RuleNodeId(toUUID(strRuleNodeId)); checkRuleNode(ruleNodeId, Operation.READ); TenantId tenantId = getCurrentUser().getTenantId(); - List events = eventService.findLatestEvents(tenantId, ruleNodeId, DataConstants.DEBUG_RULE_NODE, 2); + List events = eventService.findLatestEvents(tenantId, ruleNodeId, DataConstants.DEBUG_RULE_NODE, 2); JsonNode result = null; if (events != null) { - for (Event event : events) { + for (EventInfo event : events) { JsonNode body = event.getBody(); if (body.has("type") && body.get("type").asText().equals("IN")) { result = body; diff --git a/application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java b/application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java index 0e928b7dd6..d4d9456a3f 100644 --- a/application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java +++ b/application/src/test/java/org/thingsboard/server/actors/stats/StatsActorTest.java @@ -18,7 +18,8 @@ package org.thingsboard.server.actors.stats; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.thingsboard.server.actors.ActorSystemContext; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; +import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.dao.event.EventService; import org.thingsboard.server.queue.discovery.TbServiceInfoProvider; diff --git a/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java index ab902b82c6..7335c10d13 100644 --- a/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/AbstractRuleEngineControllerTest.java @@ -20,7 +20,7 @@ import com.fasterxml.jackson.databind.JsonNode; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.test.context.TestPropertySource; import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; @@ -60,20 +60,20 @@ public abstract class AbstractRuleEngineControllerTest extends AbstractControlle return doGet("/api/ruleChain/metadata/" + ruleChainId.getId().toString(), RuleChainMetaData.class); } - protected PageData getDebugEvents(TenantId tenantId, EntityId entityId, int limit) throws Exception { + protected PageData getDebugEvents(TenantId tenantId, EntityId entityId, int limit) throws Exception { return getEvents(tenantId, entityId, DataConstants.DEBUG_RULE_NODE, limit); } - protected PageData getEvents(TenantId tenantId, EntityId entityId, String eventType, int limit) throws Exception { + protected PageData getEvents(TenantId tenantId, EntityId entityId, String eventType, int limit) throws Exception { TimePageLink pageLink = new TimePageLink(limit); return doGetTypedWithTimePageLink("/api/events/{entityType}/{entityId}/{eventType}?tenantId={tenantId}&", - new TypeReference>() { + new TypeReference>() { }, pageLink, entityId.getEntityType(), entityId.getId(), eventType, tenantId.getId()); } - protected JsonNode getMetadata(Event outEvent) { + protected JsonNode getMetadata(EventInfo outEvent) { String metaDataStr = outEvent.getBody().get("metadata").asText(); try { return mapper.readTree(metaDataStr); @@ -82,7 +82,7 @@ public abstract class AbstractRuleEngineControllerTest extends AbstractControlle } } - protected Predicate filterByCustomEvent() { + protected Predicate filterByCustomEvent() { return event -> event.getBody().get("msgType").textValue().equals("CUSTOM"); } diff --git a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java index 9d96fe9d9c..24c94adba6 100644 --- a/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/rules/flow/AbstractRuleEngineFlowIntegrationTest.java @@ -30,9 +30,10 @@ 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.EventInfo; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.User; +import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.page.PageData; @@ -177,15 +178,15 @@ public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRule 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()); + PageData eventsPage = getDebugEvents(savedTenant.getId(), ruleChain.getFirstRuleNodeId(), 1000); + List events = eventsPage.getData().stream().filter(filterByCustomEvent()).collect(Collectors.toList()); Assert.assertEquals(2, events.size()); - Event inEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); + EventInfo inEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); Assert.assertEquals(ruleChain.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(); + EventInfo outEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.OUT)).findFirst().get(); Assert.assertEquals(ruleChain.getFirstRuleNodeId(), outEvent.getEntityId()); Assert.assertEquals(device.getId().getId().toString(), outEvent.getBody().get("entityId").asText()); @@ -299,16 +300,16 @@ public abstract class AbstractRuleEngineFlowIntegrationTest extends AbstractRule Mockito.verify(tbMsgCallback, Mockito.timeout(10000)).onSuccess(); - PageData eventsPage = getDebugEvents(savedTenant.getId(), rootRuleChain.getFirstRuleNodeId(), 1000); - List events = eventsPage.getData().stream().filter(filterByCustomEvent()).collect(Collectors.toList()); + PageData eventsPage = getDebugEvents(savedTenant.getId(), rootRuleChain.getFirstRuleNodeId(), 1000); + List events = eventsPage.getData().stream().filter(filterByCustomEvent()).collect(Collectors.toList()); Assert.assertEquals(2, events.size()); - Event inEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); + EventInfo inEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); Assert.assertEquals(rootRuleChain.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(); + EventInfo outEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.OUT)).findFirst().get(); Assert.assertEquals(rootRuleChain.getFirstRuleNodeId(), outEvent.getEntityId()); Assert.assertEquals(device.getId().getId().toString(), outEvent.getBody().get("entityId").asText()); 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 5d86e6a2a5..c840cf4c9a 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 @@ -20,18 +20,14 @@ import org.awaitility.Awaitility; import org.junit.After; import org.junit.Assert; import org.junit.Before; -import org.junit.ClassRule; import org.junit.Test; import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.test.context.support.TestPropertySourceUtils; -import org.testcontainers.containers.GenericContainer; 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.EventInfo; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.rule.RuleChain; @@ -109,11 +105,11 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac Assert.assertNotNull(ruleChainFinal.getFirstRuleNodeId()); //TODO find out why RULE_NODE update event did not appear all the time - List rcEvents = Awaitility.await("Rule Node started successfully") + List rcEvents = Awaitility.await("Rule Node started successfully") .pollInterval(10, MILLISECONDS) .atMost(TIMEOUT, TimeUnit.SECONDS) .until(() -> { - List debugEvents = getEvents(tenantId, ruleChainFinal.getFirstRuleNodeId(), DataConstants.LC_EVENT, 1000) + List debugEvents = getEvents(tenantId, ruleChainFinal.getFirstRuleNodeId(), DataConstants.LC_EVENT, 1000) .getData().stream().filter(e -> { var body = e.getBody(); return body.has("event") && body.get("event").asText().equals("STARTED") @@ -145,11 +141,11 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac 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") + List events = Awaitility.await("get debug by custom event") .pollInterval(10, MILLISECONDS) .atMost(TIMEOUT, TimeUnit.SECONDS) .until(() -> { - List debugEvents = getDebugEvents(tenantId, ruleChainFinal.getFirstRuleNodeId(), 1000) + 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)); @@ -158,11 +154,11 @@ public abstract class AbstractRuleEngineLifecycleIntegrationTest extends Abstrac x -> x.size() == 2); log.warn("asserting.."); - Event inEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); + EventInfo inEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.IN)).findFirst().get(); 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(); + EventInfo outEvent = events.stream().filter(e -> e.getBody().get("type").asText().equals(DataConstants.OUT)).findFirst().get(); Assert.assertEquals(ruleChainFinal.getFirstRuleNodeId(), outEvent.getEntityId()); Assert.assertEquals(device.getId().getId().toString(), outEvent.getBody().get("entityId").asText()); diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java index 83624df8e3..0d157e755d 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/event/EventService.java @@ -16,7 +16,8 @@ package org.thingsboard.server.dao.event; import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; +import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.event.EventFilter; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -30,15 +31,15 @@ public interface EventService { ListenableFuture saveAsync(Event event); - Optional findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid); + Optional findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid); - PageData findEvents(TenantId tenantId, EntityId entityId, TimePageLink pageLink); + PageData findEvents(TenantId tenantId, EntityId entityId, TimePageLink pageLink); - PageData findEvents(TenantId tenantId, EntityId entityId, String eventType, TimePageLink pageLink); + PageData findEvents(TenantId tenantId, EntityId entityId, String eventType, TimePageLink pageLink); - List findLatestEvents(TenantId tenantId, EntityId entityId, String eventType, int limit); + List findLatestEvents(TenantId tenantId, EntityId entityId, String eventType, int limit); - PageData findEventsByFilter(TenantId tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink); + PageData findEventsByFilter(TenantId tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink); void removeEvents(TenantId tenantId, EntityId entityId); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/Event.java b/common/data/src/main/java/org/thingsboard/server/common/data/EventInfo.java similarity index 93% rename from common/data/src/main/java/org/thingsboard/server/common/data/Event.java rename to common/data/src/main/java/org/thingsboard/server/common/data/EventInfo.java index 29ab33d075..f6c4ee69fb 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/Event.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/EventInfo.java @@ -28,7 +28,7 @@ import org.thingsboard.server.common.data.id.TenantId; */ @Data @ApiModel -public class Event extends BaseData { +public class EventInfo extends BaseData { @ApiModelProperty(position = 1, value = "JSON object with Tenant Id.", accessMode = ApiModelProperty.AccessMode.READ_ONLY) private TenantId tenantId; @@ -41,15 +41,15 @@ public class Event extends BaseData { @ApiModelProperty(position = 5, value = "Event body.", dataType = "com.fasterxml.jackson.databind.JsonNode") private transient JsonNode body; - public Event() { + public EventInfo() { super(); } - public Event(EventId id) { + public EventInfo(EventId id) { super(id); } - public Event(Event event) { + public EventInfo(EventInfo event) { super(event); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/ErrorEvent.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/ErrorEvent.java new file mode 100644 index 0000000000..8e06177ca6 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/ErrorEvent.java @@ -0,0 +1,48 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed 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.thingsboard.server.common.data.event; + +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.Setter; +import lombok.ToString; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; + +@ToString +@EqualsAndHashCode(callSuper = true) +public class ErrorEvent extends Event { + + private static final long serialVersionUID = 960461434033192571L; + + @Builder + private ErrorEvent(TenantId tenantId, EntityId entityId, String serviceId, String method, String error) { + super(tenantId, entityId, serviceId); + this.method = method; + this.error = error; + } + + @Getter @Setter + private String method; + @Getter @Setter + private String error; + + @Override + public EventType getType() { + return EventType.ERROR; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/Event.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/Event.java new file mode 100644 index 0000000000..9c013e1486 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/Event.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed 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.thingsboard.server.common.data.event; + +import lombok.Data; +import lombok.EqualsAndHashCode; +import org.thingsboard.server.common.data.BaseData; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.EventId; +import org.thingsboard.server.common.data.id.TenantId; + +@Data +@EqualsAndHashCode(callSuper = true) +public abstract class Event extends BaseData { + + private final TenantId tenantId; + private final EntityId entityId; + private final String serviceId; + + public Event(TenantId tenantId, EntityId entityId, String serviceId) { + super(); + this.tenantId = tenantId != null ? tenantId : TenantId.SYS_TENANT_ID; + this.entityId = entityId; + this.serviceId = serviceId; + } + + public abstract EventType getType(); + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java index 2645e7cce3..8d58158521 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/EventType.java @@ -15,6 +15,16 @@ */ package org.thingsboard.server.common.data.event; +import lombok.Getter; + public enum EventType { - ERROR, LC_EVENT, STATS, DEBUG_RULE_NODE, DEBUG_RULE_CHAIN + ERROR("error_event"), LC_EVENT("lc_event"), STATS("stats_event"), DEBUG_RULE_NODE("rule_node_debug_event"), DEBUG_RULE_CHAIN("rule_chain_debug_event"); + + @Getter + private final String table; + + EventType(String table) { + this.table = table; + } + } \ No newline at end of file diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/LifecycleEvent.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/LifecycleEvent.java new file mode 100644 index 0000000000..217955af96 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/LifecycleEvent.java @@ -0,0 +1,51 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed 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.thingsboard.server.common.data.event; + +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.Setter; +import lombok.ToString; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; + +@ToString +@EqualsAndHashCode(callSuper = true) +public class LifecycleEvent extends Event { + + private static final long serialVersionUID = -3247420461850911549L; + + @Builder + private LifecycleEvent(TenantId tenantId, EntityId entityId, String serviceId, String lcEventType, boolean success, String error) { + super(tenantId, entityId, serviceId); + this.lcEventType = lcEventType; + this.success = success; + this.error = error; + } + + @Getter + private final String lcEventType; + @Getter + private final boolean success; + @Getter @Setter + private String error; + + @Override + public EventType getType() { + return EventType.LC_EVENT; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/RuleChainDebugEvent.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/RuleChainDebugEvent.java new file mode 100644 index 0000000000..ac8c16c2e5 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/RuleChainDebugEvent.java @@ -0,0 +1,50 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed 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.thingsboard.server.common.data.event; + +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.Setter; +import lombok.ToString; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; + +import java.util.UUID; + +@ToString +@EqualsAndHashCode(callSuper = true) +public class RuleChainDebugEvent extends Event { + + private static final long serialVersionUID = -386392236201116767L; + + @Builder + private RuleChainDebugEvent(TenantId tenantId, EntityId entityId, String serviceId, String message, String error) { + super(tenantId, entityId, serviceId); + this.message = message; + this.error = error; + } + + @Getter @Setter + private String message; + @Getter @Setter + private String error; + + @Override + public EventType getType() { + return EventType.DEBUG_RULE_CHAIN; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/RuleNodeDebugEvent.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/RuleNodeDebugEvent.java new file mode 100644 index 0000000000..f9280d2d6b --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/RuleNodeDebugEvent.java @@ -0,0 +1,75 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed 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.thingsboard.server.common.data.event; + +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.Setter; +import lombok.ToString; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; + +import java.util.UUID; + +@ToString +@EqualsAndHashCode(callSuper = true) +public class RuleNodeDebugEvent extends Event { + + private static final long serialVersionUID = -6575797430064573984L; + + @Builder + private RuleNodeDebugEvent(TenantId tenantId, EntityId entityId, String serviceId, + String eventType, EntityId eventEntity, UUID msgId, + String msgType, String dataType, String relationType, + String data, String metadata, String error) { + super(tenantId, entityId, serviceId); + this.eventType = eventType; + this.eventEntity = eventEntity; + this.msgId = msgId; + this.msgType = msgType; + this.dataType = dataType; + this.relationType = relationType; + this.data = data; + this.metadata = metadata; + this.error = error; + } + + @Getter + private final String eventType; + @Getter + private final EntityId eventEntity; + @Getter + private final UUID msgId; + @Getter + private final String msgType; + @Getter + private final String dataType; + @Getter + private final String relationType; + @Getter @Setter + private String data; + @Getter @Setter + private String metadata; + @Getter @Setter + private String error; + + //TODO: rename the enum constant + @Override + public EventType getType() { + return EventType.DEBUG_RULE_NODE; + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/event/StatisticsEvent.java b/common/data/src/main/java/org/thingsboard/server/common/data/event/StatisticsEvent.java new file mode 100644 index 0000000000..b0194762de --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/event/StatisticsEvent.java @@ -0,0 +1,47 @@ +/** + * Copyright © 2016-2022 The Thingsboard Authors + * + * Licensed 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.thingsboard.server.common.data.event; + +import lombok.Builder; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.ToString; +import org.thingsboard.server.common.data.id.EntityId; +import org.thingsboard.server.common.data.id.TenantId; + +@ToString +@EqualsAndHashCode(callSuper = true) +public class StatisticsEvent extends Event { + + private static final long serialVersionUID = 6683733979448910631L; + + @Builder + private StatisticsEvent(TenantId tenantId, EntityId entityId, String serviceId, long messagesProcessed, long errorsOccurred) { + super(tenantId, entityId, serviceId); + this.messagesProcessed = messagesProcessed; + this.errorsOccurred = errorsOccurred; + } + + @Getter + private final long messagesProcessed; + @Getter + private final long errorsOccurred; + + @Override + public EventType getType() { + return EventType.STATS; + } +} diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java index 3937bbd9c3..d7f0885bdc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java @@ -15,15 +15,19 @@ */ package org.thingsboard.server.dao.event; -import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; +import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.event.ErrorEvent; +import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.event.EventFilter; +import org.thingsboard.server.common.data.event.LifecycleEvent; +import org.thingsboard.server.common.data.event.RuleChainDebugEvent; +import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.IdBased; import org.thingsboard.server.common.data.id.TenantId; @@ -34,6 +38,8 @@ import org.thingsboard.server.dao.service.DataValidator; import java.util.List; import java.util.Optional; +import java.util.function.BiConsumer; +import java.util.function.Function; import java.util.stream.Collectors; @Service @@ -57,18 +63,41 @@ public class BaseEventService implements EventService { } private void checkAndTruncateDebugEvent(Event event) { - if (event.getType().startsWith("DEBUG") && event.getBody() != null && event.getBody().has("data")) { - String dataStr = event.getBody().get("data").asText(); - int length = dataStr.length(); + switch (event.getType()) { + case DEBUG_RULE_NODE: + RuleNodeDebugEvent rnEvent = (RuleNodeDebugEvent) event; + truncateField(rnEvent, RuleNodeDebugEvent::getData, RuleNodeDebugEvent::setData); + truncateField(rnEvent, RuleNodeDebugEvent::getMetadata, RuleNodeDebugEvent::setMetadata); + truncateField(rnEvent, RuleNodeDebugEvent::getError, RuleNodeDebugEvent::setError); + break; + case DEBUG_RULE_CHAIN: + RuleChainDebugEvent rcEvent = (RuleChainDebugEvent) event; + truncateField(rcEvent, RuleChainDebugEvent::getMessage, RuleChainDebugEvent::setMessage); + truncateField(rcEvent, RuleChainDebugEvent::getError, RuleChainDebugEvent::setError); + break; + case LC_EVENT: + LifecycleEvent lcEvent = (LifecycleEvent) event; + truncateField(lcEvent, LifecycleEvent::getError, LifecycleEvent::setError); + break; + case ERROR: + ErrorEvent eEvent = (ErrorEvent) event; + truncateField(eEvent, ErrorEvent::getError, ErrorEvent::setError); + break; + } + } + + private void truncateField(T event, Function getter, BiConsumer setter) { + var str = getter.apply(event); + if (StringUtils.isNotEmpty(str)) { + var length = str.length(); if (length > maxDebugEventSymbols) { - ((ObjectNode) event.getBody()).put("data", dataStr.substring(0, maxDebugEventSymbols) + "...[truncated " + (length - maxDebugEventSymbols) + " symbols]"); - log.trace("[{}] Event was truncated: {}", event.getId(), dataStr); + setter.accept(event, str.substring(0, maxDebugEventSymbols) + "...[truncated " + (length - maxDebugEventSymbols) + " symbols]"); } } } @Override - public Optional findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) { + public Optional findEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) { if (tenantId == null) { throw new DataValidationException("Tenant id should be specified!."); } @@ -81,27 +110,27 @@ public class BaseEventService implements EventService { if (StringUtils.isEmpty(eventUid)) { throw new DataValidationException("Event uid should be specified!."); } - Event event = eventDao.findEvent(tenantId.getId(), entityId, eventType, eventUid); + EventInfo event = eventDao.findEvent(tenantId.getId(), entityId, eventType, eventUid); return event != null ? Optional.of(event) : Optional.empty(); } @Override - public PageData findEvents(TenantId tenantId, EntityId entityId, TimePageLink pageLink) { + public PageData findEvents(TenantId tenantId, EntityId entityId, TimePageLink pageLink) { return eventDao.findEvents(tenantId.getId(), entityId, pageLink); } @Override - public PageData findEvents(TenantId tenantId, EntityId entityId, String eventType, TimePageLink pageLink) { + public PageData findEvents(TenantId tenantId, EntityId entityId, String eventType, TimePageLink pageLink) { return eventDao.findEvents(tenantId.getId(), entityId, eventType, pageLink); } @Override - public List findLatestEvents(TenantId tenantId, EntityId entityId, String eventType, int limit) { + public List findLatestEvents(TenantId tenantId, EntityId entityId, String eventType, int limit) { return eventDao.findLatestEvents(tenantId.getId(), entityId, eventType, limit); } @Override - public PageData findEventsByFilter(TenantId tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink) { + public PageData findEventsByFilter(TenantId tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink) { return eventDao.findEventByFilter(tenantId.getId(), entityId, eventFilter, pageLink); } @@ -113,7 +142,7 @@ public class BaseEventService implements EventService { @Override public void removeEvents(TenantId tenantId, EntityId entityId, EventFilter eventFilter, Long startTime, Long endTime) { TimePageLink eventsPageLink = new TimePageLink(1000, 0, null, null, startTime, endTime); - PageData eventsPageData; + PageData eventsPageData; do { if (eventFilter == null) { eventsPageData = findEvents(tenantId, entityId, eventsPageLink); diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java index 86eb0d91ec..30815bd15e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java @@ -16,7 +16,8 @@ package org.thingsboard.server.dao.event; import com.google.common.util.concurrent.ListenableFuture; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; +import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.event.EventFilter; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.page.PageData; @@ -48,7 +49,7 @@ public interface EventDao extends Dao { * @param eventUid the eventUid * @return the event */ - Event findEvent(UUID tenantId, EntityId entityId, String eventType, String eventUid); + EventInfo findEvent(UUID tenantId, EntityId entityId, String eventType, String eventUid); /** * Find events by tenantId, entityId and pageLink. @@ -58,7 +59,7 @@ public interface EventDao extends Dao { * @param pageLink the pageLink * @return the event list */ - PageData findEvents(UUID tenantId, EntityId entityId, TimePageLink pageLink); + PageData findEvents(UUID tenantId, EntityId entityId, TimePageLink pageLink); /** * Find events by tenantId, entityId, eventType and pageLink. @@ -69,9 +70,9 @@ public interface EventDao extends Dao { * @param pageLink the pageLink * @return the event list */ - PageData findEvents(UUID tenantId, EntityId entityId, String eventType, TimePageLink pageLink); + PageData findEvents(UUID tenantId, EntityId entityId, String eventType, TimePageLink pageLink); - PageData findEventByFilter(UUID tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink); + PageData findEventByFilter(UUID tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink); /** * Find latest events by tenantId, entityId and eventType. @@ -82,7 +83,7 @@ public interface EventDao extends Dao { * @param limit the limit * @return the event list */ - List findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit); + List findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit); /** * Executes stored procedure to cleanup old events. Uses separate ttl for debug and other events. diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/EventEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/EventEntity.java index b452ec5bda..bf08333e7b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/EventEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/EventEntity.java @@ -22,7 +22,8 @@ import lombok.NoArgsConstructor; import org.hibernate.annotations.Type; import org.hibernate.annotations.TypeDef; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; +import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.EventId; import org.thingsboard.server.common.data.id.TenantId; @@ -78,7 +79,7 @@ public class EventEntity extends BaseSqlEntity implements BaseEntity implements BaseEntity { @Override protected void validateDataImpl(TenantId tenantId, Event event) { + if (event.getTenantId() == null) { + throw new DataValidationException("Tenant id should be specified!."); + } if (event.getEntityId() == null) { throw new DataValidationException("Entity id should be specified!."); } - if (StringUtils.isEmpty(event.getType())) { - throw new DataValidationException("Event type should be specified!."); - } - if (event.getBody() == null) { - throw new DataValidationException("Event body should be specified!."); + if (StringUtils.isEmpty(event.getServiceId())) { + throw new DataValidationException("Service id should be specified!."); } } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java index 6234db58da..e69b8578df 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.dao.sql.event; +import org.jetbrains.annotations.NotNull; +import org.postgresql.core.Oid; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.jdbc.core.BatchPreparedStatementSetter; @@ -24,12 +26,24 @@ import org.springframework.transaction.TransactionStatus; import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.support.TransactionCallbackWithoutResult; import org.springframework.transaction.support.TransactionTemplate; -import org.thingsboard.server.dao.model.sql.EventEntity; +import org.thingsboard.server.common.data.event.ErrorEvent; +import org.thingsboard.server.common.data.event.Event; +import org.thingsboard.server.common.data.event.EventType; +import org.thingsboard.server.common.data.event.LifecycleEvent; +import org.thingsboard.server.common.data.event.RuleChainDebugEvent; +import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; +import org.thingsboard.server.common.data.event.StatisticsEvent; +import javax.annotation.PostConstruct; import java.sql.PreparedStatement; import java.sql.SQLException; +import java.sql.Types; import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; import java.util.regex.Pattern; +import java.util.stream.Collectors; @Repository @Transactional @@ -39,10 +53,7 @@ public class EventInsertRepository { private static final String EMPTY_STR = ""; - private static final String INSERT = - "INSERT INTO event (id, created_time, body, entity_id, entity_type, event_type, event_uid, tenant_id, ts) " + - "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) " + - "ON CONFLICT DO NOTHING;"; + private final Map insertStmtMap = new ConcurrentHashMap<>(); @Autowired protected JdbcTemplate jdbcTemplate; @@ -53,34 +64,172 @@ public class EventInsertRepository { @Value("${sql.remove_null_chars:true}") private boolean removeNullChars; - protected void save(List entities) { + @PostConstruct + public void init() { + insertStmtMap.put(EventType.ERROR, "INSERT INTO " + EventType.ERROR.getTable() + + " (id, tenant_id, ts, entity_id, service_id, e_method, e_error) " + + "VALUES (?, ?, ?, ?, ?, ?, ?) ON CONFLICT DO NOTHING;"); + insertStmtMap.put(EventType.LC_EVENT, "INSERT INTO " + EventType.LC_EVENT.getTable() + + " (id, tenant_id, ts, entity_id, service_id, e_type, e_success, e_error) " + + "VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT DO NOTHING;"); + insertStmtMap.put(EventType.STATS, "INSERT INTO " + EventType.STATS.getTable() + + " (id, tenant_id, ts, entity_id, service_id, e_messages_processed, e_errors_occured) " + + "VALUES (?, ?, ?, ?, ?, ?, ?) ON CONFLICT DO NOTHING;"); + insertStmtMap.put(EventType.DEBUG_RULE_NODE, "INSERT INTO " + EventType.DEBUG_RULE_NODE.getTable() + + " (id, tenant_id, ts, entity_id, service_id, e_type, e_entity_id, e_entity_type, e_msg_id, e_msg_type, e_data_type, e_relation_type, e_data, e_metadata, e_error) " + + "VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT DO NOTHING;"); + insertStmtMap.put(EventType.DEBUG_RULE_CHAIN, "INSERT INTO " + EventType.DEBUG_RULE_CHAIN.getTable() + + " (id, tenant_id, ts, entity_id, service_id, e_message, e_error) " + + "VALUES (?, ?, ?, ?, ?, ?, ?) ON CONFLICT DO NOTHING;"); + } + + protected void save(List entities) { + Map> eventsByType = entities.stream().collect(Collectors.groupingBy(Event::getType, Collectors.toList())); transactionTemplate.execute(new TransactionCallbackWithoutResult() { @Override protected void doInTransactionWithoutResult(TransactionStatus status) { - jdbcTemplate.batchUpdate(INSERT, new BatchPreparedStatementSetter() { - @Override - public void setValues(PreparedStatement ps, int i) throws SQLException { - EventEntity event = entities.get(i); - ps.setObject(1, event.getId()); - ps.setLong(2, event.getCreatedTime()); - ps.setString(3, replaceNullChars(event.getBody().toString())); - ps.setObject(4, event.getEntityId()); - ps.setString(5, event.getEntityType().name()); - ps.setString(6, event.getEventType()); - ps.setString(7, event.getEventUid()); - ps.setObject(8, event.getTenantId()); - ps.setLong(9, event.getTs()); - } - - @Override - public int getBatchSize() { - return entities.size(); - } - }); + for (var entry : eventsByType.entrySet()) { + jdbcTemplate.batchUpdate(insertStmtMap.get(entry.getKey()), getStatementSetter(entry.getKey(), entry.getValue())); + } } }); } + private BatchPreparedStatementSetter getStatementSetter(EventType eventType, List events) { + switch (eventType) { + case ERROR: + return getErrorEventSetter(events); + case LC_EVENT: + return getLcEventSetter(events); + case STATS: + return getStatsEventSetter(events); + case DEBUG_RULE_NODE: + return getRuleNodeEventSetter(events); + case DEBUG_RULE_CHAIN: + return getRuleChainEventSetter(events); + default: + throw new RuntimeException(eventType + " support is not implemented!"); + } + } + + private BatchPreparedStatementSetter getErrorEventSetter(List events) { + return new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement ps, int i) throws SQLException { + ErrorEvent event = (ErrorEvent) events.get(i); + setCommonEventFields(ps, event); + safePutString(ps, 6, event.getMethod()); + safePutString(ps, 7, event.getError()); + } + + @Override + public int getBatchSize() { + return events.size(); + } + }; + } + + private BatchPreparedStatementSetter getLcEventSetter(List events) { + return new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement ps, int i) throws SQLException { + LifecycleEvent event = (LifecycleEvent) events.get(i); + setCommonEventFields(ps, event); + safePutString(ps, 6, event.getLcEventType()); + ps.setBoolean(7, event.isSuccess()); + safePutString(ps, 8, event.getError()); + } + + @Override + public int getBatchSize() { + return events.size(); + } + }; + } + + private BatchPreparedStatementSetter getStatsEventSetter(List events) { + return new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement ps, int i) throws SQLException { + StatisticsEvent event = (StatisticsEvent) events.get(i); + setCommonEventFields(ps, event); + ps.setLong(6, event.getMessagesProcessed()); + ps.setLong(7, event.getErrorsOccurred()); + } + + @Override + public int getBatchSize() { + return events.size(); + } + }; + } + + private BatchPreparedStatementSetter getRuleNodeEventSetter(List events) { + return new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement ps, int i) throws SQLException { + RuleNodeDebugEvent event = (RuleNodeDebugEvent) events.get(i); + setCommonEventFields(ps, event); + safePutString(ps, 6, event.getEventType()); + safePutUUID(ps, 7, event.getEntityId() != null ? event.getEntityId().getId() : null); + safePutString(ps, 8, event.getEntityId() != null ? event.getEntityId().getEntityType().name() : null); + safePutUUID(ps, 9, event.getMsgId()); + safePutString(ps, 10, event.getMsgType()); + safePutString(ps, 11, event.getDataType()); + safePutString(ps, 12, event.getRelationType()); + safePutString(ps, 13, event.getData()); + safePutString(ps, 14, event.getMetadata()); + safePutString(ps, 15, event.getError()); + } + + @Override + public int getBatchSize() { + return events.size(); + } + }; + } + + private BatchPreparedStatementSetter getRuleChainEventSetter(List events) { + return new BatchPreparedStatementSetter() { + @Override + public void setValues(PreparedStatement ps, int i) throws SQLException { + RuleChainDebugEvent event = (RuleChainDebugEvent) events.get(i); + setCommonEventFields(ps, event); + safePutString(ps, 6, event.getMessage()); + safePutString(ps, 7, event.getError()); + } + + @Override + public int getBatchSize() { + return events.size(); + } + }; + } + + void safePutString(PreparedStatement ps, int parameterIdx, String value) throws SQLException { + if (value != null) { + ps.setString(parameterIdx, replaceNullChars(value)); + } else { + ps.setNull(parameterIdx, Types.VARCHAR); + } + } + + void safePutUUID(PreparedStatement ps, int parameterIdx, UUID value) throws SQLException { + if (value != null) { + ps.setObject(parameterIdx, value); + } else { + ps.setNull(parameterIdx, Types.OTHER); + } + } + + private void setCommonEventFields(PreparedStatement ps, Event event) throws SQLException { + ps.setObject(1, event.getId().getId()); + ps.setObject(2, event.getTenantId().getId()); + ps.setLong(3, event.getCreatedTime()); + ps.setObject(4, event.getEntityId().getId()); + ps.setString(5, event.getServiceId()); + } + private String replaceNullChars(String strValue) { if (removeNullChars && strValue != null) { return PATTERN_THREAD_LOCAL.get().matcher(strValue).replaceAll(EMPTY_STR); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java index 2b409006d2..deb3c87402 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java @@ -19,15 +19,19 @@ import com.datastax.oss.driver.api.core.uuid.Uuids; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; +import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.springframework.dao.DataIntegrityViolationException; import org.springframework.data.domain.PageRequest; import org.springframework.data.jpa.repository.JpaRepository; import org.springframework.stereotype.Component; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.event.DebugEvent; import org.thingsboard.server.common.data.event.ErrorEventFilter; +import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.event.EventFilter; +import org.thingsboard.server.common.data.event.EventType; import org.thingsboard.server.common.data.event.LifeCycleEventFilter; import org.thingsboard.server.common.data.event.StatisticsEventFilter; import org.thingsboard.server.common.data.id.EntityId; @@ -42,17 +46,27 @@ import org.thingsboard.server.dao.sql.JpaAbstractDao; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; +import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository; +import org.thingsboard.server.dao.timeseries.SqlPartition; +import org.thingsboard.server.dao.timeseries.SqlTsPartitionDate; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneOffset; +import java.time.ZonedDateTime; +import java.time.format.DateTimeFormatter; import java.util.Comparator; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.ReentrantLock; import java.util.function.Function; -import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; - /** * Created by Valerii Sosliuk on 5/3/2017. */ @@ -60,7 +74,12 @@ import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; @Component public class JpaBaseEventDao extends JpaAbstractDao implements EventDao { - private final UUID systemTenantId = NULL_UUID; + private static final long PARTITION_DURATION = TimeUnit.HOURS.toMillis(1); + private final Map> partitionsByEventType = new ConcurrentHashMap<>(); + private static final ReentrantLock partitionCreationLock = new ReentrantLock(); + + @Autowired + private SqlPartitioningRepository partitioningRepository; @Autowired private EventRepository eventRepository; @@ -102,10 +121,13 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen @Value("${sql.batch_sort:false}") private boolean batchSortEnabled; - private TbSqlBlockingQueueWrapper queue; + private TbSqlBlockingQueueWrapper queue; @PostConstruct private void init() { + for (EventType eventType : EventType.values()) { + partitionsByEventType.put(eventType, new ConcurrentHashMap<>()); + } TbSqlBlockingQueueParams params = TbSqlBlockingQueueParams.builder() .logName("Events") .batchSize(batchSize) @@ -114,11 +136,9 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen .statsNamePrefix("events") .batchSortEnabled(batchSortEnabled) .build(); - Function hashcodeFunction = entity -> entity.getEntityId().hashCode(); + Function hashcodeFunction = entity -> Objects.hash(super.hashCode(), entity.getTenantId(), entity.getEntityId()); queue = new TbSqlBlockingQueueWrapper<>(params, hashcodeFunction, batchThreads, statsFactory); - queue.init(logExecutor, v -> eventInsertRepository.save(v), - Comparator.comparing((EventEntity eventEntity) -> eventEntity.getTs()) - ); + queue.init(logExecutor, v -> eventInsertRepository.save(v), Comparator.comparing(Event::getCreatedTime)); } @PreDestroy @@ -143,68 +163,80 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen event.setCreatedTime(System.currentTimeMillis()); } } - if (StringUtils.isEmpty(event.getUid())) { - event.setUid(event.getId().toString()); - } - - return save(new EventEntity(event)); + savePartitionIfNotExist(event); + return queue.add(event); } - private ListenableFuture save(EventEntity entity) { - log.debug("Save event [{}] ", entity); - if (entity.getTenantId() == null) { - log.trace("Save system event with predefined id {}", systemTenantId); - entity.setTenantId(systemTenantId); - } - if (entity.getUuid() == null) { - entity.setUuid(Uuids.timeBased()); - } - if (StringUtils.isEmpty(entity.getEventUid())) { - entity.setEventUid(entity.getUuid().toString()); + private void savePartitionIfNotExist(Event event) { + var partitionsMap = partitionsByEventType.get(event.getType()); + long partitionStartTs = event.getCreatedTime() - (event.getCreatedTime() % PARTITION_DURATION); + if (partitionsMap.get(partitionStartTs) == null) { + long partitionEndTs = partitionStartTs + PARTITION_DURATION; + savePartition(partitionsMap, new SqlPartition(event.getType().getTable(), partitionStartTs, partitionEndTs, Long.toString(partitionStartTs))); } - return addToQueue(entity); } - private ListenableFuture addToQueue(EventEntity entity) { - return queue.add(entity); + private void savePartition(Map partitionsMap, SqlPartition sqlPartition) { + if (!partitionsMap.containsKey(sqlPartition.getStart())) { + partitionCreationLock.lock(); + try { + log.trace("Saving partition: {}", sqlPartition); + partitioningRepository.save(sqlPartition); + log.trace("Adding partition to map: {}", sqlPartition); + partitionsMap.put(sqlPartition.getStart(), sqlPartition); + } catch (DataIntegrityViolationException ex) { + log.trace("Error occurred during partition save:", ex); + if (ex.getCause() instanceof ConstraintViolationException) { + log.warn("Saving partition [{}] rejected. Event data will save to the DEFAULT partition.", sqlPartition.getPartitionDate()); + partitionsMap.put(sqlPartition.getStart(), sqlPartition); + } else { + throw new RuntimeException(ex); + } + } finally { + partitionCreationLock.unlock(); + } + } } @Override - public Event findEvent(UUID tenantId, EntityId entityId, String eventType, String eventUid) { - return DaoUtil.getData(eventRepository.findByTenantIdAndEntityTypeAndEntityIdAndEventTypeAndEventUid( - tenantId, entityId.getEntityType(), entityId.getId(), eventType, eventUid)); + public EventInfo findEvent(UUID tenantId, EntityId entityId, String eventType, String eventUid) { + return null; +// return DaoUtil.getData(eventRepository.findByTenantIdAndEntityTypeAndEntityIdAndEventTypeAndEventUid( +// tenantId, entityId.getEntityType(), entityId.getId(), eventType, eventUid)); } @Override - public PageData findEvents(UUID tenantId, EntityId entityId, TimePageLink pageLink) { - return DaoUtil.toPageData( - eventRepository - .findEventsByTenantIdAndEntityId( - tenantId, - entityId.getEntityType(), - entityId.getId(), - Objects.toString(pageLink.getTextSearch(), ""), - pageLink.getStartTime(), - pageLink.getEndTime(), - DaoUtil.toPageable(pageLink))); + public PageData findEvents(UUID tenantId, EntityId entityId, TimePageLink pageLink) { + return null; +// return DaoUtil.toPageData( +// eventRepository +// .findEventsByTenantIdAndEntityId( +// tenantId, +// entityId.getEntityType(), +// entityId.getId(), +// Objects.toString(pageLink.getTextSearch(), ""), +// pageLink.getStartTime(), +// pageLink.getEndTime(), +// DaoUtil.toPageable(pageLink))); } @Override - public PageData findEvents(UUID tenantId, EntityId entityId, String eventType, TimePageLink pageLink) { - return DaoUtil.toPageData( - eventRepository - .findEventsByTenantIdAndEntityIdAndEventType( - tenantId, - entityId.getEntityType(), - entityId.getId(), - eventType, - pageLink.getStartTime(), - pageLink.getEndTime(), - DaoUtil.toPageable(pageLink))); + public PageData findEvents(UUID tenantId, EntityId entityId, String eventType, TimePageLink pageLink) { + return null; +// return DaoUtil.toPageData( +// eventRepository +// .findEventsByTenantIdAndEntityIdAndEventType( +// tenantId, +// entityId.getEntityType(), +// entityId.getId(), +// eventType, +// pageLink.getStartTime(), +// pageLink.getEndTime(), +// DaoUtil.toPageable(pageLink))); } @Override - public PageData findEventByFilter(UUID tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink) { + public PageData findEventByFilter(UUID tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink) { if (eventFilter.hasFilterForJsonBody()) { switch (eventFilter.getEventType()) { case DEBUG_RULE_NODE: @@ -224,86 +256,91 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen } } - private PageData findEventByFilter(UUID tenantId, EntityId entityId, DebugEvent eventFilter, TimePageLink pageLink) { - return DaoUtil.toPageData( - eventRepository.findDebugRuleNodeEvents( - tenantId, - entityId.getId(), - entityId.getEntityType().name(), - eventFilter.getEventType().name(), - notNull(pageLink.getStartTime()), - notNull(pageLink.getEndTime()), - eventFilter.getMsgDirectionType(), - eventFilter.getServer(), - eventFilter.getEntityName(), - eventFilter.getRelationType(), - eventFilter.getEntityId(), - eventFilter.getMsgType(), - eventFilter.isError(), - eventFilter.getErrorStr(), - eventFilter.getDataSearch(), - eventFilter.getMetadataSearch(), - DaoUtil.toPageable(pageLink))); + private PageData findEventByFilter(UUID tenantId, EntityId entityId, DebugEvent eventFilter, TimePageLink pageLink) { + return null; +// return DaoUtil.toPageData( +// eventRepository.findDebugRuleNodeEvents( +// tenantId, +// entityId.getId(), +// entityId.getEntityType().name(), +// eventFilter.getEventType().name(), +// notNull(pageLink.getStartTime()), +// notNull(pageLink.getEndTime()), +// eventFilter.getMsgDirectionType(), +// eventFilter.getServer(), +// eventFilter.getEntityName(), +// eventFilter.getRelationType(), +// eventFilter.getEntityId(), +// eventFilter.getMsgType(), +// eventFilter.isError(), +// eventFilter.getErrorStr(), +// eventFilter.getDataSearch(), +// eventFilter.getMetadataSearch(), +// DaoUtil.toPageable(pageLink))); } - private PageData findEventByFilter(UUID tenantId, EntityId entityId, ErrorEventFilter eventFilter, TimePageLink pageLink) { - return DaoUtil.toPageData( - eventRepository.findErrorEvents( - tenantId, - entityId.getId(), - entityId.getEntityType().name(), - notNull(pageLink.getStartTime()), - notNull(pageLink.getEndTime()), - eventFilter.getServer(), - eventFilter.getMethod(), - eventFilter.getErrorStr(), - DaoUtil.toPageable(pageLink)) - ); + private PageData findEventByFilter(UUID tenantId, EntityId entityId, ErrorEventFilter eventFilter, TimePageLink pageLink) { + return null; +// return DaoUtil.toPageData( +// eventRepository.findErrorEvents( +// tenantId, +// entityId.getId(), +// entityId.getEntityType().name(), +// notNull(pageLink.getStartTime()), +// notNull(pageLink.getEndTime()), +// eventFilter.getServer(), +// eventFilter.getMethod(), +// eventFilter.getErrorStr(), +// DaoUtil.toPageable(pageLink)) +// ); } - private PageData findEventByFilter(UUID tenantId, EntityId entityId, LifeCycleEventFilter eventFilter, TimePageLink pageLink) { - boolean statusFilterEnabled = !StringUtils.isEmpty(eventFilter.getStatus()); - boolean statusFilter = statusFilterEnabled && eventFilter.getStatus().equalsIgnoreCase("Success"); - return DaoUtil.toPageData( - eventRepository.findLifeCycleEvents( - tenantId, - entityId.getId(), - entityId.getEntityType().name(), - notNull(pageLink.getStartTime()), - notNull(pageLink.getEndTime()), - eventFilter.getServer(), - eventFilter.getEvent(), - statusFilterEnabled, - statusFilter, - eventFilter.getErrorStr(), - DaoUtil.toPageable(pageLink)) - ); + private PageData findEventByFilter(UUID tenantId, EntityId entityId, LifeCycleEventFilter eventFilter, TimePageLink pageLink) { + return null; +// boolean statusFilterEnabled = !StringUtils.isEmpty(eventFilter.getStatus()); +// boolean statusFilter = statusFilterEnabled && eventFilter.getStatus().equalsIgnoreCase("Success"); +// return DaoUtil.toPageData( +// eventRepository.findLifeCycleEvents( +// tenantId, +// entityId.getId(), +// entityId.getEntityType().name(), +// notNull(pageLink.getStartTime()), +// notNull(pageLink.getEndTime()), +// eventFilter.getServer(), +// eventFilter.getEvent(), +// statusFilterEnabled, +// statusFilter, +// eventFilter.getErrorStr(), +// DaoUtil.toPageable(pageLink)) +// ); } - private PageData findEventByFilter(UUID tenantId, EntityId entityId, StatisticsEventFilter eventFilter, TimePageLink pageLink) { - return DaoUtil.toPageData( - eventRepository.findStatisticsEvents( - tenantId, - entityId.getId(), - entityId.getEntityType().name(), - notNull(pageLink.getStartTime()), - notNull(pageLink.getEndTime()), - eventFilter.getServer(), - notNull(eventFilter.getMessagesProcessed()), - notNull(eventFilter.getErrorsOccurred()), - DaoUtil.toPageable(pageLink)) - ); + private PageData findEventByFilter(UUID tenantId, EntityId entityId, StatisticsEventFilter eventFilter, TimePageLink pageLink) { + return null; +// return DaoUtil.toPageData( +// eventRepository.findStatisticsEvents( +// tenantId, +// entityId.getId(), +// entityId.getEntityType().name(), +// notNull(pageLink.getStartTime()), +// notNull(pageLink.getEndTime()), +// eventFilter.getServer(), +// notNull(eventFilter.getMessagesProcessed()), +// notNull(eventFilter.getErrorsOccurred()), +// DaoUtil.toPageable(pageLink)) +// ); } @Override - public List findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit) { - List latest = eventRepository.findLatestByTenantIdAndEntityTypeAndEntityIdAndEventType( - tenantId, - entityId.getEntityType(), - entityId.getId(), - eventType, - PageRequest.of(0, limit)); - return DaoUtil.convertDataList(latest); + public List findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit) { + return null; +// List latest = eventRepository.findLatestByTenantIdAndEntityTypeAndEntityIdAndEventType( +// tenantId, +// entityId.getEntityType(), +// entityId.getId(), +// eventType, +// PageRequest.of(0, limit)); +// return DaoUtil.convertDataList(latest); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java index 899a195538..805e5e4f4f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java @@ -23,7 +23,6 @@ import org.thingsboard.server.dao.util.SqlTsDao; import javax.persistence.EntityManager; import javax.persistence.PersistenceContext; -@SqlTsDao @Repository @Transactional public class SqlPartitioningRepository { @@ -32,8 +31,7 @@ public class SqlPartitioningRepository { private EntityManager entityManager; public void save(SqlPartition partition) { - entityManager.createNativeQuery(partition.getQuery()) - .executeUpdate(); + entityManager.createNativeQuery(partition.getQuery()).executeUpdate(); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/sql/JpaSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/sql/JpaSqlTimeseriesDao.java index ce173e046b..64d1669eae 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/sql/JpaSqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/sql/JpaSqlTimeseriesDao.java @@ -132,7 +132,7 @@ public class JpaSqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDao long partitionEndTs = toMills(localDateTimeEnd); ZonedDateTime zonedDateTime = localDateTimeStart.atZone(ZoneOffset.UTC); String partitionDate = zonedDateTime.format(DateTimeFormatter.ofPattern(tsFormat.getPattern())); - savePartition(new SqlPartition(partitionStartTs, partitionEndTs, partitionDate)); + savePartition(new SqlPartition(SqlPartition.TS_KV, partitionStartTs, partitionEndTs, partitionDate)); } } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/SqlPartition.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/SqlPartition.java index 5c47689b6d..dcdb335784 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/SqlPartition.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/SqlPartition.java @@ -20,21 +20,21 @@ import lombok.Data; @Data public class SqlPartition { - private static final String TABLE_REGEX = "ts_kv_"; + public static final String TS_KV = "ts_kv"; private long start; private long end; private String partitionDate; private String query; - public SqlPartition(long start, long end, String partitionDate) { + public SqlPartition(String table, long start, long end, String partitionDate) { this.start = start; this.end = end; this.partitionDate = partitionDate; - this.query = createStatement(start, end, partitionDate); + this.query = createStatement(table, start, end, partitionDate); } - private String createStatement(long start, long end, String partitionDate) { - return "CREATE TABLE IF NOT EXISTS " + TABLE_REGEX + partitionDate + " PARTITION OF ts_kv FOR VALUES FROM (" + start + ") TO (" + end + ")"; + private String createStatement(String table, long start, long end, String partitionDate) { + return "CREATE TABLE IF NOT EXISTS " + table + "_" + partitionDate + " PARTITION OF " + table + " FOR VALUES FROM (" + start + ") TO (" + end + ")"; } } \ No newline at end of file diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index a9b25a66f7..9494e4bdd6 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -323,18 +323,64 @@ CREATE TABLE IF NOT EXISTS device_credentials ( CONSTRAINT device_credentials_device_id_unq_key UNIQUE (device_id) ); -CREATE TABLE IF NOT EXISTS event ( - id uuid NOT NULL CONSTRAINT event_pkey PRIMARY KEY, - created_time bigint NOT NULL, - body varchar(10000000), - entity_id uuid, - entity_type varchar(255), - event_type varchar(255), - event_uid varchar(255), - tenant_id uuid, +CREATE TABLE IF NOT EXISTS rule_node_debug_event ( + id uuid NOT NULL CONSTRAINT rule_node_debug_event_pkey PRIMARY KEY, + tenant_id uuid NOT NULL , ts bigint NOT NULL, - CONSTRAINT event_unq_key UNIQUE (tenant_id, entity_type, entity_id, event_type, event_uid) -); + entity_id uuid NOT NULL, + service_id varchar, + e_type varchar, + e_entity_id uuid, + e_entity_type varchar, + e_msg_id uuid, + e_msg_type varchar, + e_data_type varchar, + e_relation_type varchar, + e_data varchar NOT NULL, + e_metadata varchar NOT NULL, + e_error varchar NOT NULL +) PARTITION BY RANGE (ts); + +CREATE TABLE IF NOT EXISTS rule_chain_debug_event ( + id uuid NOT NULL CONSTRAINT rule_chain_debug_event_pkey PRIMARY KEY, + tenant_id uuid NOT NULL, + ts bigint NOT NULL, + entity_id uuid NOT NULL, + service_id varchar NOT NULL, + e_message varchar, + e_error varchar +) PARTITION BY RANGE (ts); + +CREATE TABLE IF NOT EXISTS stats_event ( + id uuid NOT NULL CONSTRAINT stats_event_pkey PRIMARY KEY, + tenant_id uuid NOT NULL, + ts bigint NOT NULL, + entity_id uuid NOT NULL, + service_id varchar NOT NULL, + e_messages_processed bigint NOT NULL, + e_errors_occured bigint NOT NULL +) PARTITION BY RANGE (ts); + +CREATE TABLE IF NOT EXISTS lc_event ( + id uuid NOT NULL CONSTRAINT lc_event_pkey PRIMARY KEY, + tenant_id uuid NOT NULL, + ts bigint NOT NULL, + entity_id uuid NOT NULL, + service_id varchar NOT NULL, + e_type varchar NOT NULL, + e_success boolean NOT NULL, + e_error varchar +) PARTITION BY RANGE (ts); + +CREATE TABLE IF NOT EXISTS error_event ( + id uuid NOT NULL CONSTRAINT error_event_pkey PRIMARY KEY, + tenant_id uuid NOT NULL, + ts bigint NOT NULL, + entity_id uuid NOT NULL, + service_id varchar NOT NULL, + e_method varchar NOT NULL, + e_error varchar +) PARTITION BY RANGE (ts); CREATE TABLE IF NOT EXISTS relation ( from_id uuid, diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java index 964e6229f1..1276f46621 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java @@ -28,17 +28,20 @@ import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringRunner; import org.springframework.test.context.support.AnnotationConfigContextLoader; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileType; import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.OtaPackage; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileConfiguration; import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.DeviceProfileData; import org.thingsboard.server.common.data.edge.Edge; +import org.thingsboard.server.common.data.event.Event; +import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.HasId; @@ -185,17 +188,16 @@ public abstract class AbstractServiceTest { } - protected Event generateEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) throws IOException { + protected RuleNodeDebugEvent generateEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) throws IOException { if (tenantId == null) { tenantId = TenantId.fromUUID(Uuids.timeBased()); } - Event event = new Event(); - event.setTenantId(tenantId); - event.setEntityId(entityId); - event.setType(eventType); - event.setUid(eventUid); - event.setBody(readFromResource("TestJsonData.json")); - return event; + return RuleNodeDebugEvent.builder() + .tenantId(tenantId) + .entityId(entityId) + .serviceId("server A") + .data(JacksonUtil.toString(readFromResource("TestJsonData.json"))) + .build(); } // // private ComponentDescriptor getOrCreateDescriptor(ComponentScope scope, ComponentType type, String clazz, String configurationDescriptorResource) throws IOException { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java index f3f6dd8ff8..ffd15ff5ba 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/event/BaseEventServiceTest.java @@ -16,12 +16,13 @@ package org.thingsboard.server.dao.service.event; import com.datastax.oss.driver.api.core.uuid.Uuids; -import org.apache.commons.lang3.time.DateFormatUtils; import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; +import org.thingsboard.server.common.data.event.Event; +import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; @@ -33,9 +34,6 @@ import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.service.AbstractServiceTest; import java.text.ParseException; -import java.time.LocalDateTime; -import java.time.Month; -import java.time.ZoneOffset; import java.util.Optional; import static org.apache.commons.lang3.time.DateFormatUtils.ISO_DATETIME_TIME_ZONE_FORMAT; @@ -59,14 +57,15 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { @Test public void saveEvent() throws Exception { DeviceId devId = new DeviceId(Uuids.timeBased()); - Event event = generateEvent(null, devId, "ALARM", Uuids.timeBased().toString()); + RuleNodeDebugEvent event = generateEvent(null, devId, "ALARM", Uuids.timeBased().toString()); eventService.saveAsync(event).get(); - Optional loaded = eventService.findEvent(event.getTenantId(), event.getEntityId(), event.getType(), event.getUid()); - Assert.assertTrue(loaded.isPresent()); - Assert.assertNotNull(loaded.get()); - Assert.assertEquals(event.getEntityId(), loaded.get().getEntityId()); - Assert.assertEquals(event.getType(), loaded.get().getType()); - Assert.assertEquals(event.getBody(), loaded.get().getBody()); + throw new RuntimeException("fix me!"); +// Optional loaded = eventService.findEvent(event.getTenantId(), event.getEntityId(), event.getType(), event.getUid()); +// Assert.assertTrue(loaded.isPresent()); +// Assert.assertNotNull(loaded.get()); +// Assert.assertEquals(event.getEntityId(), loaded.get().getEntityId()); +// Assert.assertEquals(event.getType(), loaded.get().getType()); +// Assert.assertEquals(event.getBody(), loaded.get().getBody()); } @Test @@ -74,14 +73,14 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { CustomerId customerId = new CustomerId(Uuids.timeBased()); TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); saveEventWithProvidedTime(timeBeforeStartTime, customerId, tenantId); - Event savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); - Event savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); - Event savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); + EventInfo savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); + EventInfo savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); + EventInfo savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); saveEventWithProvidedTime(timeAfterEndTime, customerId, tenantId); TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime"), startTime, endTime); - PageData events = eventService.findEvents(tenantId, customerId, DataConstants.STATS, + PageData events = eventService.findEvents(tenantId, customerId, DataConstants.STATS, timePageLink); Assert.assertNotNull(events.getData()); @@ -105,14 +104,14 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { CustomerId customerId = new CustomerId(Uuids.timeBased()); TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); saveEventWithProvidedTime(timeBeforeStartTime, customerId, tenantId); - Event savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); - Event savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); - Event savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); + EventInfo savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); + EventInfo savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); + EventInfo savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); saveEventWithProvidedTime(timeAfterEndTime, customerId, tenantId); TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime", SortOrder.Direction.DESC), startTime, endTime); - PageData events = eventService.findEvents(tenantId, customerId, DataConstants.STATS, + PageData events = eventService.findEvents(tenantId, customerId, DataConstants.STATS, timePageLink); Assert.assertNotNull(events.getData()); @@ -131,10 +130,11 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1, timeBeforeStartTime - 1, timeAfterEndTime + 1); } - private Event saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws Exception { - Event event = generateEvent(tenantId, entityId, DataConstants.STATS, null); - event.setId(new EventId(Uuids.startOf(time))); - eventService.saveAsync(event).get(); - return event; + private EventInfo saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws Exception { + throw new RuntimeException("fix me!"); +// EventInfo event = generateEvent(tenantId, entityId, DataConstants.STATS, null); +// event.setId(new EventId(Uuids.startOf(time))); +// eventService.saveAsync(event).get(); +// return event; } } diff --git a/dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java index eb5f004714..4381aab18f 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDaoTest.java @@ -23,7 +23,7 @@ import org.junit.After; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; import org.thingsboard.server.common.data.DataConstants; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EventId; @@ -35,9 +35,7 @@ import org.thingsboard.server.dao.event.EventDao; import java.io.IOException; import java.util.List; -import java.util.Optional; import java.util.UUID; -import java.util.concurrent.ExecutionException; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -58,128 +56,129 @@ public class JpaBaseEventDaoTest extends AbstractJpaDaoTest { @After public void deleteEvents() { - List events = eventDao.find(TenantId.fromUUID(tenantId)); - for (Event event : events) { - eventDao.removeById(TenantId.fromUUID(tenantId), event.getUuidId()); - } + throw new RuntimeException("fix me!"); +// List events = eventDao.find(TenantId.fromUUID(tenantId)); +// for (EventInfo event : events) { +// eventDao.removeById(TenantId.fromUUID(tenantId), event.getUuidId()); +// } } - @Test - public void findEvent() { - UUID entityId = Uuids.timeBased(); - Event savedEvent = eventDao.save(TenantId.fromUUID(tenantId), getEvent(entityId, tenantId, entityId)); - Event foundEvent = eventDao.findEvent(tenantId, new DeviceId(entityId), DataConstants.STATS, savedEvent.getUid()); - assertNotNull("Event expected to be not null", foundEvent); - assertEquals(savedEvent.getId(), foundEvent.getId()); - } - - @Test - public void findEventsByEntityIdAndPageLink() throws Exception { - UUID entityId1 = Uuids.timeBased(); - UUID entityId2 = Uuids.timeBased(); - long startTime = System.currentTimeMillis(); - long endTime = createEventsTwoEntities(tenantId, entityId1, entityId2, 20); - - TimePageLink pageLink1 = new TimePageLink(30); - PageData events1 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink1); - assertEquals(10, events1.getData().size()); - - TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null); - PageData events2 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink2); - assertEquals(10, events2.getData().size()); - - TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime); - PageData events3 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink3); - assertEquals(10, events3.getData().size()); - - TimePageLink pageLink4 = new TimePageLink(5, 0, "", null, startTime, endTime); - PageData events4 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4); - assertEquals(5, events4.getData().size()); - - pageLink4 = pageLink4.nextPageLink(); - PageData events5 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4); - assertEquals(5, events5.getData().size()); - - pageLink4 = pageLink4.nextPageLink(); - PageData events6 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4); - assertEquals(0, events6.getData().size()); - - } - - @Test - public void findEventsByEntityIdAndEventTypeAndPageLink() throws Exception { - UUID entityId1 = Uuids.timeBased(); - UUID entityId2 = Uuids.timeBased(); - long startTime = System.currentTimeMillis(); - long endTime = createEventsTwoEntitiesTwoTypes(tenantId, entityId1, entityId2, 20); - - TimePageLink pageLink1 = new TimePageLink(30); - PageData events1 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink1); - assertEquals(5, events1.getData().size()); - - TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null); - PageData events2 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink2); - assertEquals(5, events2.getData().size()); - - TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime); - PageData events3 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink3); - assertEquals(5, events3.getData().size()); - - TimePageLink pageLink4 = new TimePageLink(4, 0, "", null, startTime, endTime); - PageData events4 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink4); - assertEquals(4, events4.getData().size()); - - pageLink4 = pageLink4.nextPageLink(); - PageData events5 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink4); - assertEquals(1, events5.getData().size()); - } - - private long createEventsTwoEntitiesTwoTypes(UUID tenantId, UUID entityId1, UUID entityId2, int count) throws Exception { - for (int i = 0; i < count / 2; i++) { - String type = i % 2 == 0 ? STATS : ALARM; - UUID eventId1 = Uuids.timeBased(); - Event event1 = getEvent(eventId1, tenantId, entityId1, type); - eventDao.saveAsync(event1).get(); - UUID eventId2 = Uuids.timeBased(); - Event event2 = getEvent(eventId2, tenantId, entityId2, type); - eventDao.saveAsync(event2).get(); - } - return System.currentTimeMillis(); - } - - private long createEventsTwoEntities(UUID tenantId, UUID entityId1, UUID entityId2, int count) throws Exception { - for (int i = 0; i < count / 2; i++) { - UUID eventId1 = Uuids.timeBased(); - Event event1 = getEvent(eventId1, tenantId, entityId1); - eventDao.saveAsync(event1).get(); - UUID eventId2 = Uuids.timeBased(); - Event event2 = getEvent(eventId2, tenantId, entityId2); - eventDao.saveAsync(event2).get(); - } - return System.currentTimeMillis(); - } - - private Event getEvent(UUID eventId, UUID tenantId, UUID entityId, String type) { - Event event = getEvent(eventId, tenantId, entityId); - event.setType(type); - return event; - } - - private Event getEvent(UUID eventId, UUID tenantId, UUID entityId) { - Event event = new Event(); - event.setId(new EventId(eventId)); - event.setTenantId(TenantId.fromUUID(tenantId)); - EntityId deviceId = new DeviceId(entityId); - event.setEntityId(deviceId); - event.setUid(event.getId().getId().toString()); - event.setType(STATS); - ObjectMapper mapper = new ObjectMapper(); - try { - JsonNode jsonNode = mapper.readTree("{\"key\":\"value\"}"); - event.setBody(jsonNode); - } catch (IOException e) { - log.error(e.getMessage(), e); - } - return event; - } +// @Test +// public void findEvent() { +// UUID entityId = Uuids.timeBased(); +// EventInfo savedEvent = eventDao.save(TenantId.fromUUID(tenantId), getEvent(entityId, tenantId, entityId)); +// EventInfo foundEvent = eventDao.findEvent(tenantId, new DeviceId(entityId), DataConstants.STATS, savedEvent.getUid()); +// assertNotNull("Event expected to be not null", foundEvent); +// assertEquals(savedEvent.getId(), foundEvent.getId()); +// } +// +// @Test +// public void findEventsByEntityIdAndPageLink() throws Exception { +// UUID entityId1 = Uuids.timeBased(); +// UUID entityId2 = Uuids.timeBased(); +// long startTime = System.currentTimeMillis(); +// long endTime = createEventsTwoEntities(tenantId, entityId1, entityId2, 20); +// +// TimePageLink pageLink1 = new TimePageLink(30); +// PageData events1 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink1); +// assertEquals(10, events1.getData().size()); +// +// TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null); +// PageData events2 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink2); +// assertEquals(10, events2.getData().size()); +// +// TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime); +// PageData events3 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink3); +// assertEquals(10, events3.getData().size()); +// +// TimePageLink pageLink4 = new TimePageLink(5, 0, "", null, startTime, endTime); +// PageData events4 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4); +// assertEquals(5, events4.getData().size()); +// +// pageLink4 = pageLink4.nextPageLink(); +// PageData events5 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4); +// assertEquals(5, events5.getData().size()); +// +// pageLink4 = pageLink4.nextPageLink(); +// PageData events6 = eventDao.findEvents(tenantId, new DeviceId(entityId1), pageLink4); +// assertEquals(0, events6.getData().size()); +// +// } +// +// @Test +// public void findEventsByEntityIdAndEventTypeAndPageLink() throws Exception { +// UUID entityId1 = Uuids.timeBased(); +// UUID entityId2 = Uuids.timeBased(); +// long startTime = System.currentTimeMillis(); +// long endTime = createEventsTwoEntitiesTwoTypes(tenantId, entityId1, entityId2, 20); +// +// TimePageLink pageLink1 = new TimePageLink(30); +// PageData events1 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink1); +// assertEquals(5, events1.getData().size()); +// +// TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null); +// PageData events2 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink2); +// assertEquals(5, events2.getData().size()); +// +// TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime); +// PageData events3 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink3); +// assertEquals(5, events3.getData().size()); +// +// TimePageLink pageLink4 = new TimePageLink(4, 0, "", null, startTime, endTime); +// PageData events4 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink4); +// assertEquals(4, events4.getData().size()); +// +// pageLink4 = pageLink4.nextPageLink(); +// PageData events5 = eventDao.findEvents(tenantId, new DeviceId(entityId1), ALARM, pageLink4); +// assertEquals(1, events5.getData().size()); +// } +// +// private long createEventsTwoEntitiesTwoTypes(UUID tenantId, UUID entityId1, UUID entityId2, int count) throws Exception { +// for (int i = 0; i < count / 2; i++) { +// String type = i % 2 == 0 ? STATS : ALARM; +// UUID eventId1 = Uuids.timeBased(); +// EventInfo event1 = getEvent(eventId1, tenantId, entityId1, type); +// eventDao.saveAsync(event1).get(); +// UUID eventId2 = Uuids.timeBased(); +// EventInfo event2 = getEvent(eventId2, tenantId, entityId2, type); +// eventDao.saveAsync(event2).get(); +// } +// return System.currentTimeMillis(); +// } +// +// private long createEventsTwoEntities(UUID tenantId, UUID entityId1, UUID entityId2, int count) throws Exception { +// for (int i = 0; i < count / 2; i++) { +// UUID eventId1 = Uuids.timeBased(); +// EventInfo event1 = getEvent(eventId1, tenantId, entityId1); +// eventDao.saveAsync(event1).get(); +// UUID eventId2 = Uuids.timeBased(); +// EventInfo event2 = getEvent(eventId2, tenantId, entityId2); +// eventDao.saveAsync(event2).get(); +// } +// return System.currentTimeMillis(); +// } +// +// private EventInfo getEvent(UUID eventId, UUID tenantId, UUID entityId, String type) { +// EventInfo event = getEvent(eventId, tenantId, entityId); +// event.setType(type); +// return event; +// } +// +// private EventInfo getEvent(UUID eventId, UUID tenantId, UUID entityId) { +// EventInfo event = new EventInfo(); +// event.setId(new EventId(eventId)); +// event.setTenantId(TenantId.fromUUID(tenantId)); +// EntityId deviceId = new DeviceId(entityId); +// event.setEntityId(deviceId); +// event.setUid(event.getId().getId().toString()); +// event.setType(STATS); +// ObjectMapper mapper = new ObjectMapper(); +// try { +// JsonNode jsonNode = mapper.readTree("{\"key\":\"value\"}"); +// event.setBody(jsonNode); +// } catch (IOException e) { +// log.error(e.getMessage(), e); +// } +// return event; +// } } diff --git a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java index a0f1d5dfa1..f2f446e534 100644 --- a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java +++ b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java @@ -54,7 +54,7 @@ import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.EntityViewInfo; -import org.thingsboard.server.common.data.Event; +import org.thingsboard.server.common.data.EventInfo; import org.thingsboard.server.common.data.OtaPackage; import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.SaveDeviceWithCredentialsRequest; @@ -1811,7 +1811,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { } } - public PageData getEvents(EntityId entityId, String eventType, TenantId tenantId, TimePageLink pageLink) { + public PageData getEvents(EntityId entityId, String eventType, TenantId tenantId, TimePageLink pageLink) { Map params = new HashMap<>(); params.put("entityType", entityId.getEntityType().name()); params.put("entityId", entityId.getId().toString()); @@ -1823,12 +1823,12 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/events/{entityType}/{entityId}/{eventType}?tenantId={tenantId}&" + getTimeUrlParams(pageLink), HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference>() { + new ParameterizedTypeReference>() { }, params).getBody(); } - public PageData getEvents(EntityId entityId, TenantId tenantId, TimePageLink pageLink) { + public PageData getEvents(EntityId entityId, TenantId tenantId, TimePageLink pageLink) { Map params = new HashMap<>(); params.put("entityType", entityId.getEntityType().name()); params.put("entityId", entityId.getId().toString()); @@ -1839,7 +1839,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { baseURL + "/api/events/{entityType}/{entityId}?tenantId={tenantId}&" + getTimeUrlParams(pageLink), HttpMethod.GET, HttpEntity.EMPTY, - new ParameterizedTypeReference>() { + new ParameterizedTypeReference>() { }, params).getBody(); }