From a38edcd62b54ce96bb121f5e9cdb10d1d3d501a6 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Wed, 27 Jul 2022 18:56:09 +0300 Subject: [PATCH] Upgrade script for the events migration --- .../main/data/upgrade/3.4.0/schema_update.sql | 215 ++++++++++++++++++ .../install/ThingsboardInstallService.java | 3 + .../install/SqlDatabaseUpgradeService.java | 12 + .../service/ttl/EventsCleanUpService.java | 8 +- .../server/dao/event/EventService.java | 3 +- .../server/common/data/event/ErrorEvent.java | 15 +- .../common/data/event/RuleNodeDebugEvent.java | 3 +- .../common/data/event/StatisticsEvent.java | 6 +- .../server/dao/event/BaseEventService.java | 4 +- .../server/dao/event/EventDao.java | 7 +- .../dao/sql/event/ErrorEventRepository.java | 4 +- .../server/dao/sql/event/JpaBaseEventDao.java | 21 +- .../sql/event/LifecycleEventRepository.java | 5 +- .../event/RuleChainDebugEventRepository.java | 4 +- .../event/RuleNodeDebugEventRepository.java | 5 +- .../sql/event/StatisticsEventRepository.java | 4 +- .../dao/service/AbstractServiceTest.java | 2 +- .../service/event/BaseEventServiceTest.java | 77 +++---- .../dao/sql/event/JpaBaseEventDaoTest.java | 214 +++++++---------- 19 files changed, 409 insertions(+), 203 deletions(-) create mode 100644 application/src/main/data/upgrade/3.4.0/schema_update.sql diff --git a/application/src/main/data/upgrade/3.4.0/schema_update.sql b/application/src/main/data/upgrade/3.4.0/schema_update.sql new file mode 100644 index 0000000000..2b4b5af42d --- /dev/null +++ b/application/src/main/data/upgrade/3.4.0/schema_update.sql @@ -0,0 +1,215 @@ +-- +-- 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. +-- + +CREATE TABLE IF NOT EXISTS rule_node_debug_event ( + id uuid NOT NULL, + tenant_id uuid NOT NULL , + ts bigint NOT NULL, + 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, + e_metadata varchar, + e_error varchar +) PARTITION BY RANGE (ts); + +CREATE TABLE IF NOT EXISTS rule_chain_debug_event ( + id uuid NOT NULL, + 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, + 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_occurred bigint NOT NULL +) PARTITION BY RANGE (ts); + +CREATE TABLE IF NOT EXISTS lc_event ( + id uuid NOT NULL, + 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, + 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 INDEX IF NOT EXISTS idx_rule_node_debug_event_main + ON rule_node_debug_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95); + +CREATE INDEX IF NOT EXISTS idx_rule_chain_debug_event_main + ON rule_chain_debug_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95); + +CREATE INDEX IF NOT EXISTS idx_stats_event_main + ON stats_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95); + +CREATE INDEX IF NOT EXISTS idx_lc_event_main + ON lc_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95); + +CREATE INDEX IF NOT EXISTS idx_error_event_main + ON error_event (tenant_id ASC, entity_id ASC, ts DESC NULLS LAST) WITH (FILLFACTOR=95); + + +-- Useful to migrate old events to the new table structure; +CREATE OR REPLACE PROCEDURE migrate_regular_events(IN start_ts_in_ms bigint, IN partition_size_in_hours int) + LANGUAGE plpgsql AS +$$ +DECLARE + partition_size_in_ms bigint; + p record; + table_name varchar; +BEGIN + partition_size_in_ms = partition_size_in_hours * 3600 * 1000; + + FOR p IN SELECT DISTINCT event_type as event_type, (created_time - created_time % partition_size_in_ms) as partition_ts FROM event e WHERE e.event_type in ('STATS', 'LC_EVENT', 'ERROR') and ts > start_ts_in_ms + LOOP + IF p.event_type = 'STATS' THEN + table_name := 'stats_event'; + ELSEIF p.event_type = 'LC_EVENT' THEN + table_name := 'lc_event'; + ELSEIF p.event_type = 'ERROR' THEN + table_name := 'error_event'; + END IF; + RAISE NOTICE '[%] Partition to create : [%-%]', table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms); + EXECUTE format('CREATE TABLE IF NOT EXISTS %s_%s PARTITION OF %s FOR VALUES FROM ( %s ) TO ( %s )', table_name, p.partition_ts, table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms)); + END LOOP; + + INSERT INTO stats_event + SELECT id, + tenant_id, + ts, + entity_id, + body::json ->> 'server', + (body::json ->> 'messagesProcessed')::bigint, + (body::json ->> 'errorsOccurred')::bigint + FROM event + WHERE ts > start_ts_in_ms + AND event_type = 'STATS' + ON CONFLICT DO NOTHING; + + INSERT INTO lc_event + SELECT id, + tenant_id, + ts, + entity_id, + body::json ->> 'server', + body::json ->> 'event', + (body::json ->> 'success')::boolean, + body::json ->> 'error' + FROM event + WHERE ts > start_ts_in_ms + AND event_type = 'LC_EVENT' + ON CONFLICT DO NOTHING; + + INSERT INTO error_event + SELECT id, + tenant_id, + ts, + entity_id, + body::json ->> 'server', + body::json ->> 'method', + body::json ->> 'error' + FROM event + WHERE ts > start_ts_in_ms + AND event_type = 'ERROR' + ON CONFLICT DO NOTHING; + +END +$$; + +-- Useful to migrate old debug events to the new table structure; +CREATE OR REPLACE PROCEDURE migrate_debug_events(IN start_ts_in_ms bigint, IN partition_size_in_hours int) + LANGUAGE plpgsql AS +$$ +DECLARE + partition_size_in_ms bigint; + p record; + table_name varchar; +BEGIN + partition_size_in_ms = partition_size_in_hours * 3600 * 1000; + + FOR p IN SELECT DISTINCT event_type as event_type, (created_time - created_time % partition_size_in_ms) as partition_ts FROM event e WHERE e.event_type in ('DEBUG_RULE_NODE', 'DEBUG_RULE_CHAIN') and ts > start_ts_in_ms + LOOP + IF p.event_type = 'DEBUG_RULE_NODE' THEN + table_name := 'rule_node_debug_event'; + ELSEIF p.event_type = 'DEBUG_RULE_CHAIN' THEN + table_name := 'rule_chain_debug_event'; + END IF; + RAISE NOTICE '[%] Partition to create : [%-%]', table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms); + EXECUTE format('CREATE TABLE IF NOT EXISTS %s_%s PARTITION OF %s FOR VALUES FROM ( %s ) TO ( %s )', table_name, p.partition_ts, table_name, p.partition_ts, (p.partition_ts + partition_size_in_ms)); + END LOOP; + + INSERT INTO rule_node_debug_event + SELECT id, + tenant_id, + ts, + entity_id, + body::json ->> 'server', + body::json ->> 'type', + (body::json ->> 'entityId')::uuid, + body::json ->> 'entityName', + (body::json ->> 'msgId')::uuid, + body::json ->> 'msgType', + body::json ->> 'dataType', + body::json ->> 'relationType', + body::json ->> 'data', + body::json ->> 'metadata', + body::json ->> 'error' + FROM event + WHERE ts > start_ts_in_ms + AND event_type = 'DEBUG_RULE_NODE' + ON CONFLICT DO NOTHING; + + INSERT INTO rule_chain_debug_event + SELECT id, + tenant_id, + ts, + entity_id, + body::json ->> 'server', + body::json ->> 'message', + body::json ->> 'error' + FROM event + WHERE ts > start_ts_in_ms + AND event_type = 'DEBUG_RULE_CHAIN' + ON CONFLICT DO NOTHING; +END +$$; diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index 40e15c7b0e..ea7f76e1c3 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -221,6 +221,9 @@ public class ThingsboardInstallService { log.info("Upgrading ThingsBoard from version 3.3.4 to 3.4.0 ..."); databaseEntitiesUpgradeService.upgradeDatabase("3.3.4"); dataUpdateService.updateData("3.3.4"); + case "3.4.0": + log.info("Upgrading ThingsBoard from version 3.4.0 to 3.4.1 ..."); + databaseEntitiesUpgradeService.upgradeDatabase("3.4.0"); log.info("Updating system data..."); systemDataLoaderService.updateSystemWidgets(); break; diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java index d22039683d..b8cdf491d7 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlDatabaseUpgradeService.java @@ -596,6 +596,18 @@ public class SqlDatabaseUpgradeService implements DatabaseEntitiesUpgradeService log.error("Failed updating schema!!!", e); } break; + case "3.4.0": + try (Connection conn = DriverManager.getConnection(dbUrl, dbUserName, dbPassword)) { + log.info("Updating schema ..."); + schemaUpdateFile = Paths.get(installScripts.getDataDir(), "upgrade", "3.4.0", SCHEMA_UPDATE_SQL); + loadSql(schemaUpdateFile, conn); + log.info("Updating schema settings..."); + conn.createStatement().execute("UPDATE tb_schema_settings SET schema_version = 3004001;"); + log.info("Schema updated."); + } catch (Exception e) { + log.error("Failed updating schema!!!", e); + } + break; default: throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java index 08ee57119c..3dc66b6210 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java @@ -25,7 +25,6 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import java.util.concurrent.TimeUnit; -@TbCoreComponent @Slf4j @Service public class EventsCleanUpService extends AbstractCleanUpService { @@ -39,9 +38,6 @@ public class EventsCleanUpService extends AbstractCleanUpService { @Value("${sql.ttl.events.debug_events_ttl}") private long debugTtlInSec; - @Value("${sql.ttl.events.execution_interval_ms}") - private long executionIntervalInMs; - @Value("${sql.ttl.events.enabled}") private boolean ttlTaskExecutionEnabled; @@ -54,11 +50,11 @@ public class EventsCleanUpService extends AbstractCleanUpService { @Scheduled(initialDelayString = RANDOM_DELAY_INTERVAL_MS_EXPRESSION, fixedDelayString = "${sql.ttl.events.execution_interval_ms}") public void cleanUp() { - if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) { + if (ttlTaskExecutionEnabled) { long ts = System.currentTimeMillis(); long regularEventExpTs = ttlInSec > 0 ? ts - TimeUnit.SECONDS.toMillis(ttlInSec) : 0; long debugEventExpTs = debugTtlInSec > 0 ? ts - TimeUnit.SECONDS.toMillis(debugTtlInSec) : 0; - eventService.cleanupEvents(regularEventExpTs, debugEventExpTs); + eventService.cleanupEvents(regularEventExpTs, debugEventExpTs, isSystemTenantPartitionMine()); } } 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 2b89959b56..feed84598b 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 @@ -26,7 +26,6 @@ import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.TimePageLink; import java.util.List; -import java.util.Optional; public interface EventService { @@ -42,6 +41,6 @@ public interface EventService { void removeEvents(TenantId tenantId, EntityId entityId, EventFilter eventFilter, Long startTime, Long endTime); - void cleanupEvents(long regularEventExpTs, long debugEventExpTs); + void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb); } 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 index f811a118cf..341548df59 100644 --- 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.data.event; +import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.Builder; import lombok.EqualsAndHashCode; import lombok.Getter; @@ -40,9 +41,11 @@ public class ErrorEvent extends Event { this.error = error; } - @Getter @Setter + @Getter + @Setter private String method; - @Getter @Setter + @Getter + @Setter private String error; @Override @@ -52,6 +55,12 @@ public class ErrorEvent extends Event { @Override public EventInfo toInfo(EntityType entityType) { - return null; + EventInfo eventInfo = super.toInfo(entityType); + var json = (ObjectNode) eventInfo.getBody(); + json.put("method", method); + if (error != null) { + json.put("error", error); + } + return eventInfo; } } 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 index fdca84239d..ee87d976cf 100644 --- 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 @@ -84,7 +84,8 @@ public class RuleNodeDebugEvent extends Event { var json = (ObjectNode) eventInfo.getBody(); json.put("type", eventType); if (eventEntity != null) { - json.put("entityId", eventEntity.getId().toString()).put("entityType", eventEntity.getEntityType().name()); + json.put("entityId", eventEntity.getId().toString()) + .put("entityType", eventEntity.getEntityType().name()); } if (msgId != null) { json.put("msgId", msgId.toString()); 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 index b574949cf8..0e5de8263f 100644 --- 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 @@ -15,6 +15,7 @@ */ package org.thingsboard.server.common.data.event; +import com.fasterxml.jackson.databind.node.ObjectNode; import lombok.Builder; import lombok.EqualsAndHashCode; import lombok.Getter; @@ -51,6 +52,9 @@ public class StatisticsEvent extends Event { @Override public EventInfo toInfo(EntityType entityType) { - return null; + EventInfo eventInfo = super.toInfo(entityType); + var json = (ObjectNode) eventInfo.getBody(); + json.put("messagesProcessed", messagesProcessed).put("errorsOccurred", errorsOccurred); + return eventInfo; } } 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 fad7f6b310..f07a31a05e 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 @@ -125,8 +125,8 @@ public class BaseEventService implements EventService { } @Override - public void cleanupEvents(long regularEventExpTs, long debugEventExpTs) { - eventDao.cleanupEvents(regularEventExpTs, debugEventExpTs); + public void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb) { + eventDao.cleanupEvents(regularEventExpTs, debugEventExpTs, cleanupDb); } private PageData convert(EntityType entityType, PageData pd) { 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 e09ea59872..e41668f699 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,15 +16,11 @@ package org.thingsboard.server.dao.event; import com.google.common.util.concurrent.ListenableFuture; -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.event.EventType; -import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.TimePageLink; -import org.thingsboard.server.dao.Dao; import java.util.List; import java.util.UUID; @@ -70,8 +66,9 @@ public interface EventDao { * Executes stored procedure to cleanup old events. Uses separate ttl for debug and other events. * @param regularEventExpTs the expiration time of the regular events * @param debugEventExpTs the expiration time of the debug events + * @param cleanupDb */ - void cleanupEvents(long regularEventExpTs, long debugEventExpTs); + void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb); /** * Removes all events for the specified entity and time interval diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/ErrorEventRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/ErrorEventRepository.java index c00d4c0f7d..43b4d0c91d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/ErrorEventRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/ErrorEventRepository.java @@ -35,8 +35,8 @@ import java.util.UUID; public interface ErrorEventRepository extends EventRepository, JpaRepository { @Override - @Query("SELECT e FROM ErrorEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC") - List findLatestEvents(UUID tenantId, UUID entityId, int limit); + @Query(nativeQuery = true, value = "SELECT * FROM error_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit") + List findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit); @Override @Query("SELECT e FROM ErrorEventEntity e WHERE " + 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 525150fd90..907aad0a50 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 @@ -427,14 +427,29 @@ public class JpaBaseEventDao implements EventDao { } @Override - public void cleanupEvents(long regularEventExpTs, long debugEventExpTs) { + public void cleanupEvents(long regularEventExpTs, long debugEventExpTs, boolean cleanupDb) { if (regularEventExpTs > 0) { log.info("Going to cleanup regular events with exp time: {}", regularEventExpTs); - eventCleanupRepository.cleanupEvents(regularEventExpTs, false); + if (cleanupDb) { + eventCleanupRepository.cleanupEvents(regularEventExpTs, false); + } + cleanupPartitions(regularEventExpTs, false); } if (debugEventExpTs > 0) { log.info("Going to cleanup debug events with exp time: {}", debugEventExpTs); - eventCleanupRepository.cleanupEvents(debugEventExpTs, true); + if (cleanupDb) { + eventCleanupRepository.cleanupEvents(debugEventExpTs, true); + } + cleanupPartitions(debugEventExpTs, true); + } + } + + private void cleanupPartitions(long expTime, boolean isDebug) { + for (EventType eventType : EventType.values()) { + if (eventType.isDebug() == isDebug) { + Map partitions = partitionsByEventType.get(eventType); + partitions.keySet().removeIf(startTs -> startTs + partitionConfiguration.getPartitionSizeInMs(eventType) < expTime); + } } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/LifecycleEventRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/LifecycleEventRepository.java index e4a7f40256..11215e4771 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/LifecycleEventRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/LifecycleEventRepository.java @@ -24,6 +24,7 @@ import org.springframework.data.repository.query.Param; import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.event.LifecycleEvent; import org.thingsboard.server.dao.model.sql.LifecycleEventEntity; +import org.thingsboard.server.dao.model.sql.RuleChainDebugEventEntity; import java.util.List; import java.util.UUID; @@ -31,8 +32,8 @@ import java.util.UUID; public interface LifecycleEventRepository extends EventRepository, JpaRepository { @Override - @Query("SELECT e FROM LifecycleEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC") - List findLatestEvents(UUID tenantId, UUID entityId, int limit); + @Query(nativeQuery = true, value = "SELECT * FROM lc_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit") + List findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit); @Query("SELECT e FROM LifecycleEventEntity e WHERE " + diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleChainDebugEventRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleChainDebugEventRepository.java index d01a54fa30..eb8ba01cf0 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleChainDebugEventRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleChainDebugEventRepository.java @@ -34,8 +34,8 @@ import java.util.UUID; public interface RuleChainDebugEventRepository extends EventRepository, JpaRepository { @Override - @Query("SELECT e FROM RuleChainDebugEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC") - List findLatestEvents(UUID tenantId, UUID entityId, int limit); + @Query(nativeQuery = true, value = "SELECT * FROM rule_chain_debug_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit") + List findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit); @Override @Query("SELECT e FROM RuleChainDebugEventEntity e WHERE " + diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleNodeDebugEventRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleNodeDebugEventRepository.java index 98d404c36f..3df9989674 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleNodeDebugEventRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/RuleNodeDebugEventRepository.java @@ -26,6 +26,7 @@ import org.thingsboard.server.common.data.event.ErrorEvent; import org.thingsboard.server.common.data.event.RuleNodeDebugEvent; import org.thingsboard.server.dao.model.sql.ErrorEventEntity; import org.thingsboard.server.dao.model.sql.RuleNodeDebugEventEntity; +import org.thingsboard.server.dao.model.sql.StatisticsEventEntity; import java.util.List; import java.util.UUID; @@ -34,8 +35,8 @@ import java.util.UUID; public interface RuleNodeDebugEventRepository extends EventRepository, JpaRepository { @Override - @Query("SELECT e FROM RuleNodeDebugEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC") - List findLatestEvents(UUID tenantId, UUID entityId, int limit); + @Query(nativeQuery = true, value = "SELECT * FROM rule_node_debug_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit") + List findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit); @Override @Query("SELECT e FROM RuleNodeDebugEventEntity e WHERE " + diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/StatisticsEventRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/StatisticsEventRepository.java index 67116100da..3d2c4f4516 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/StatisticsEventRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/StatisticsEventRepository.java @@ -31,8 +31,8 @@ import java.util.UUID; public interface StatisticsEventRepository extends EventRepository, JpaRepository { @Override - @Query("SELECT e FROM LifecycleEventEntity e WHERE e.tenantId = :tenantId AND e.entityId = :entityId ORDER BY e.ts DESC") - List findLatestEvents(UUID tenantId, UUID entityId, int limit); + @Query(nativeQuery = true, value = "SELECT * FROM stats_event e WHERE e.tenant_id = :tenantId AND e.entity_id = :entityId ORDER BY e.ts DESC LIMIT :limit") + List findLatestEvents(@Param("tenantId") UUID tenantId, @Param("entityId") UUID entityId, @Param("limit") int limit); @Query("SELECT e FROM StatisticsEventEntity e WHERE " + "e.tenantId = :tenantId " + 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 3c856c2c70..475f84a29a 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 @@ -188,7 +188,7 @@ public abstract class AbstractServiceTest { } - protected RuleNodeDebugEvent generateEvent(TenantId tenantId, EntityId entityId, String eventType, String eventUid) throws IOException { + protected RuleNodeDebugEvent generateEvent(TenantId tenantId, EntityId entityId) throws IOException { if (tenantId == null) { tenantId = TenantId.fromUUID(Uuids.timeBased()); } 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 4bb315e5dc..26e0df8c6d 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 @@ -19,7 +19,6 @@ import com.datastax.oss.driver.api.core.uuid.Uuids; 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.EventInfo; import org.thingsboard.server.common.data.event.Event; import org.thingsboard.server.common.data.event.EventType; @@ -35,7 +34,7 @@ import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.service.AbstractServiceTest; import java.text.ParseException; -import java.util.Optional; +import java.util.List; import static org.apache.commons.lang3.time.DateFormatUtils.ISO_DATETIME_TIME_ZONE_FORMAT; @@ -57,16 +56,14 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { @Test public void saveEvent() throws Exception { + TenantId tenantId = new TenantId(Uuids.timeBased()); DeviceId devId = new DeviceId(Uuids.timeBased()); - RuleNodeDebugEvent event = generateEvent(null, devId, "ALARM", Uuids.timeBased().toString()); + RuleNodeDebugEvent event = generateEvent(tenantId, devId); eventService.saveAsync(event).get(); - 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()); + List loaded = eventService.findLatestEvents(event.getTenantId(), devId, event.getType(), 1); + Assert.assertNotNull(loaded); + Assert.assertEquals(1, loaded.size()); + Assert.assertEquals(event.getData(), loaded.get(0).getBody().get("data").asText()); } @Test @@ -74,29 +71,29 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { CustomerId customerId = new CustomerId(Uuids.timeBased()); TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); saveEventWithProvidedTime(timeBeforeStartTime, customerId, tenantId); - EventInfo savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); - EventInfo savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); - EventInfo savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); + Event savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); + Event savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); + Event savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); saveEventWithProvidedTime(timeAfterEndTime, customerId, tenantId); - TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime"), startTime, endTime); + TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("ts"), startTime, endTime); - PageData events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink); + PageData events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink); Assert.assertNotNull(events.getData()); - Assert.assertTrue(events.getData().size() == 2); - Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent.getUuidId())); - Assert.assertTrue(events.getData().get(1).getUuidId().equals(savedEvent2.getUuidId())); + Assert.assertEquals(2, events.getData().size()); + Assert.assertEquals(savedEvent.getUuidId(), events.getData().get(0).getUuidId()); + Assert.assertEquals(savedEvent2.getUuidId(), events.getData().get(1).getUuidId()); Assert.assertTrue(events.hasNext()); - events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink.nextPageLink()); + events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink.nextPageLink()); Assert.assertNotNull(events.getData()); - Assert.assertTrue(events.getData().size() == 1); - Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent3.getUuidId())); + Assert.assertEquals(1, events.getData().size()); + Assert.assertEquals(savedEvent3.getUuidId(), events.getData().get(0).getUuidId()); Assert.assertFalse(events.hasNext()); - eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1); + eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1, true); } @Test @@ -104,36 +101,36 @@ public abstract class BaseEventServiceTest extends AbstractServiceTest { CustomerId customerId = new CustomerId(Uuids.timeBased()); TenantId tenantId = TenantId.fromUUID(Uuids.timeBased()); saveEventWithProvidedTime(timeBeforeStartTime, customerId, tenantId); - EventInfo savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); - EventInfo savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); - EventInfo savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); + Event savedEvent = saveEventWithProvidedTime(eventTime, customerId, tenantId); + Event savedEvent2 = saveEventWithProvidedTime(eventTime + 1, customerId, tenantId); + Event savedEvent3 = saveEventWithProvidedTime(eventTime + 2, customerId, tenantId); saveEventWithProvidedTime(timeAfterEndTime, customerId, tenantId); - TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime", SortOrder.Direction.DESC), startTime, endTime); + TimePageLink timePageLink = new TimePageLink(2, 0, "", new SortOrder("ts", SortOrder.Direction.DESC), startTime, endTime); - PageData events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink); + PageData events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink); Assert.assertNotNull(events.getData()); - Assert.assertTrue(events.getData().size() == 2); - Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent3.getUuidId())); - Assert.assertTrue(events.getData().get(1).getUuidId().equals(savedEvent2.getUuidId())); + Assert.assertEquals(2, events.getData().size()); + Assert.assertEquals(savedEvent3.getUuidId(), events.getData().get(0).getUuidId()); + Assert.assertEquals(savedEvent2.getUuidId(), events.getData().get(1).getUuidId()); Assert.assertTrue(events.hasNext()); - events = eventService.findEvents(tenantId, customerId, EventType.STATS, timePageLink.nextPageLink()); + events = eventService.findEvents(tenantId, customerId, EventType.DEBUG_RULE_NODE, timePageLink.nextPageLink()); Assert.assertNotNull(events.getData()); - Assert.assertTrue(events.getData().size() == 1); - Assert.assertTrue(events.getData().get(0).getUuidId().equals(savedEvent.getUuidId())); + Assert.assertEquals(1, events.getData().size()); + Assert.assertEquals(savedEvent.getUuidId(), events.getData().get(0).getUuidId()); Assert.assertFalse(events.hasNext()); - eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1); + eventService.cleanupEvents(timeBeforeStartTime - 1, timeAfterEndTime + 1, true); } - 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; + private Event saveEventWithProvidedTime(long time, EntityId entityId, TenantId tenantId) throws Exception { + RuleNodeDebugEvent event = generateEvent(tenantId, entityId); + event.setId(new EventId(Uuids.timeBased())); + event.setCreatedTime(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 16170c10ea..b2091379b4 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 @@ -17,16 +17,26 @@ package org.thingsboard.server.dao.sql.event; import com.datastax.oss.driver.api.core.uuid.Uuids; import lombok.extern.slf4j.Slf4j; -import org.junit.After; +import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.thingsboard.server.common.data.event.Event; +import org.thingsboard.server.common.data.event.EventType; +import org.thingsboard.server.common.data.event.StatisticsEvent; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.page.PageData; +import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.AbstractJpaDaoTest; import org.thingsboard.server.dao.event.EventDao; +import java.util.List; import java.util.UUID; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; -/** - * Created by Valerii Sosliuk on 5/5/2017. - */ @Slf4j public class JpaBaseEventDaoTest extends AbstractJpaDaoTest { @@ -34,131 +44,77 @@ public class JpaBaseEventDaoTest extends AbstractJpaDaoTest { private EventDao eventDao; UUID tenantId = Uuids.timeBased(); - @After - public void deleteEvents() { - 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() throws InterruptedException, ExecutionException, TimeoutException { + UUID entityId = Uuids.timeBased(); + + Event event1 = getStatsEvent(Uuids.timeBased(), tenantId, entityId); + eventDao.saveAsync(event1).get(1, TimeUnit.MINUTES); + Thread.sleep(2); + Event event2 = getStatsEvent(Uuids.timeBased(), tenantId, entityId); + eventDao.saveAsync(event2).get(1, TimeUnit.MINUTES); + + List foundEvents = eventDao.findLatestEvents(tenantId, entityId, EventType.STATS, 1); + assertNotNull("Events expected to be not null", foundEvents); + assertEquals(1, foundEvents.size()); + assertEquals(event2, foundEvents.get(0)); + } + + @Test + public void findEventsByEntityIdAndPageLink() throws Exception { + UUID entityId1 = Uuids.timeBased(); + UUID entityId2 = Uuids.timeBased(); + long startTime = System.currentTimeMillis(); + + Event event1 = getStatsEvent(Uuids.timeBased(), tenantId, entityId1); + eventDao.saveAsync(event1).get(1, TimeUnit.MINUTES); + Thread.sleep(2); + Event event2 = getStatsEvent(Uuids.timeBased(), tenantId, entityId2); + eventDao.saveAsync(event2).get(1, TimeUnit.MINUTES); + + long endTime = System.currentTimeMillis(); + + PageData events1 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, new TimePageLink(30)); + assertEquals(1, events1.getData().size()); + + PageData events2 = eventDao.findEvents(tenantId, entityId2, EventType.STATS, new TimePageLink(30)); + assertEquals(1, events2.getData().size()); + + PageData events3 = eventDao.findEvents(tenantId, Uuids.timeBased(), EventType.STATS, new TimePageLink(30)); + assertEquals(0, events3.getData().size()); + + + TimePageLink pageLink2 = new TimePageLink(30, 0, "", null, startTime, null); + PageData events12 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink2); + assertEquals(1, events12.getData().size()); + assertEquals(event1, events12.getData().get(0)); + + TimePageLink pageLink3 = new TimePageLink(30, 0, "", null, startTime, endTime); + PageData events13 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink3); + assertEquals(1, events13.getData().size()); + assertEquals(event1, events13.getData().get(0)); + + TimePageLink pageLink4 = new TimePageLink(5, 0, "", null, startTime, endTime); + PageData events14 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink4); + assertEquals(1, events14.getData().size()); + assertEquals(event1, events14.getData().get(0)); + + pageLink4 = pageLink4.nextPageLink(); + PageData events6 = eventDao.findEvents(tenantId, entityId1, EventType.STATS, pageLink4); + assertEquals(0, events6.getData().size()); + } -// @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; -// } + private Event getStatsEvent(UUID eventId, UUID tenantId, UUID entityId) { + StatisticsEvent.StatisticsEventBuilder event = StatisticsEvent.builder(); + event.id(eventId); + event.ts(System.currentTimeMillis()); + event.tenantId(new TenantId(tenantId)); + event.entityId(entityId); + event.serviceId("server A"); + event.messagesProcessed(1); + event.errorsOccurred(0); + return event.build(); + } }