diff --git a/application/src/main/data/upgrade/2.4.3/schema_update_psql_drop_partitions.sql b/application/src/main/data/upgrade/2.4.3/schema_update_psql_drop_partitions.sql index fcc5c6f232..9d336e0330 100644 --- a/application/src/main/data/upgrade/2.4.3/schema_update_psql_drop_partitions.sql +++ b/application/src/main/data/upgrade/2.4.3/schema_update_psql_drop_partitions.sql @@ -18,17 +18,18 @@ CREATE OR REPLACE PROCEDURE drop_partitions_by_max_ttl(IN partition_type varchar LANGUAGE plpgsql AS $$ DECLARE - max_tenant_ttl bigint; - max_customer_ttl bigint; - max_ttl bigint; - date timestamp; - partition_by_max_ttl_date varchar; - partition_month varchar; - partition_day varchar; - partition_year varchar; - partition varchar; - partition_to_delete varchar; - + max_tenant_ttl bigint; + max_customer_ttl bigint; + max_ttl bigint; + date timestamp; + partition_by_max_ttl_date varchar; + partition_by_max_ttl_month varchar; + partition_by_max_ttl_day varchar; + partition_by_max_ttl_year varchar; + partition varchar; + partition_year integer; + partition_month integer; + partition_day integer; BEGIN SELECT max(attribute_kv.long_v) @@ -45,53 +46,138 @@ BEGIN if max_ttl IS NOT NULL AND max_ttl > 0 THEN date := to_timestamp(EXTRACT(EPOCH FROM current_timestamp) - max_ttl); partition_by_max_ttl_date := get_partition_by_max_ttl_date(partition_type, date); + RAISE NOTICE 'Date by max ttl: %', date; RAISE NOTICE 'Partition by max ttl: %', partition_by_max_ttl_date; IF partition_by_max_ttl_date IS NOT NULL THEN CASE WHEN partition_type = 'DAYS' THEN - partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); - partition_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); - partition_day := SPLIT_PART(partition_by_max_ttl_date, '_', 5); + partition_by_max_ttl_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); + partition_by_max_ttl_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); + partition_by_max_ttl_day := SPLIT_PART(partition_by_max_ttl_date, '_', 5); WHEN partition_type = 'MONTHS' THEN - partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); - partition_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); + partition_by_max_ttl_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); + partition_by_max_ttl_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); ELSE - partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); + partition_by_max_ttl_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); END CASE; - FOR partition IN SELECT tablename - FROM pg_tables - WHERE schemaname = 'public' - AND tablename like 'ts_kv_' || '%' - AND tablename != 'ts_kv_latest' - AND tablename != 'ts_kv_dictionary' - AND tablename != 'ts_kv_indefinite' - LOOP - IF partition != partition_by_max_ttl_date THEN - IF partition_year IS NOT NULL THEN - IF SPLIT_PART(partition, '_', 3)::integer < partition_year::integer THEN - partition_to_delete := partition; - ELSE - IF partition_month IS NOT NULL THEN - IF SPLIT_PART(partition, '_', 4)::integer < partition_month::integer THEN - partition_to_delete := partition; + IF partition_by_max_ttl_year IS NULL THEN + RAISE NOTICE 'Failed to remove partitions by max ttl date due to partition_by_max_ttl_year is null!'; + ELSE + IF partition_type = 'YEARS' THEN + FOR partition IN SELECT tablename + FROM pg_tables + WHERE schemaname = 'public' + AND tablename like 'ts_kv_' || '%' + AND tablename != 'ts_kv_latest' + AND tablename != 'ts_kv_dictionary' + AND tablename != 'ts_kv_indefinite' + AND tablename != partition_by_max_ttl_date + LOOP + partition_year := SPLIT_PART(partition, '_', 3)::integer; + IF partition_year < partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + END IF; + END LOOP; + ELSE + IF partition_type = 'MONTHS' THEN + IF partition_by_max_ttl_month IS NULL THEN + RAISE NOTICE 'Failed to remove months partitions by max ttl date due to partition_by_max_ttl_month is null!'; + ELSE + FOR partition IN SELECT tablename + FROM pg_tables + WHERE schemaname = 'public' + AND tablename like 'ts_kv_' || '%' + AND tablename != 'ts_kv_latest' + AND tablename != 'ts_kv_dictionary' + AND tablename != 'ts_kv_indefinite' + AND tablename != partition_by_max_ttl_date + LOOP + partition_year := SPLIT_PART(partition, '_', 3)::integer; + IF partition_year > partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; ELSE - IF partition_day IS NOT NULL THEN - IF SPLIT_PART(partition, '_', 5)::integer < partition_day::integer THEN - partition_to_delete := partition; + IF partition_year < partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + ELSE + partition_month := SPLIT_PART(partition, '_', 4)::integer; + IF partition_year = partition_by_max_ttl_year::integer THEN + IF partition_month >= partition_by_max_ttl_month::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + END IF; END IF; END IF; END IF; + END LOOP; + END IF; + ELSE + IF partition_type = 'DAYS' THEN + IF partition_by_max_ttl_month IS NULL THEN + RAISE NOTICE 'Failed to remove days partitions by max ttl date due to partition_by_max_ttl_month is null!'; + ELSE + IF partition_by_max_ttl_day IS NULL THEN + RAISE NOTICE 'Failed to remove days partitions by max ttl date due to partition_by_max_ttl_day is null!'; + ELSE + FOR partition IN SELECT tablename + FROM pg_tables + WHERE schemaname = 'public' + AND tablename like 'ts_kv_' || '%' + AND tablename != 'ts_kv_latest' + AND tablename != 'ts_kv_dictionary' + AND tablename != 'ts_kv_indefinite' + AND tablename != partition_by_max_ttl_date + LOOP + partition_year := SPLIT_PART(partition, '_', 3)::integer; + IF partition_year > partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + IF partition_year < partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + ELSE + partition_month := SPLIT_PART(partition, '_', 4)::integer; + IF partition_month > partition_by_max_ttl_month::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + IF partition_month < partition_by_max_ttl_month::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + ELSE + partition_day := SPLIT_PART(partition, '_', 5)::integer; + IF partition_day >= partition_by_max_ttl_day::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + IF partition_day < partition_by_max_ttl_day::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + END IF; + END IF; + END IF; + END IF; + END IF; + END IF; + END LOOP; END IF; END IF; END IF; - IF partition_to_delete IS NOT NULL THEN - RAISE NOTICE 'Partition to delete by max ttl: %', partition_to_delete; - EXECUTE format('DROP TABLE IF EXISTS %I', partition_to_delete); - partition_to_delete := NULL; - deleted := deleted + 1; - END IF; END IF; - END LOOP; + END IF; + END IF; END IF; END IF; END @@ -107,8 +193,6 @@ BEGIN partition := 'ts_kv_' || to_char(date, 'yyyy') || '_' || to_char(date, 'MM'); WHEN partition_type = 'YEARS' THEN partition := 'ts_kv_' || to_char(date, 'yyyy'); - WHEN partition_type = 'INDEFINITE' THEN - partition := NULL; ELSE partition := NULL; END CASE; diff --git a/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java b/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java index cd8aa88863..689b7fd9a9 100644 --- a/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java +++ b/application/src/main/java/org/thingsboard/server/service/resource/DefaultTbResourceService.java @@ -86,7 +86,10 @@ public class DefaultTbResourceService implements TbResourceService { } else { throw new DataValidationException(String.format("Could not parse the XML of objectModel with name %s", resource.getSearchText())); } - } catch (InvalidDDFFileException | IOException e) { + } catch (InvalidDDFFileException e) { + log.error("Failed to parse file {}", resource.getFileName(), e); + throw new DataValidationException("Failed to parse file " + resource.getFileName()); + } catch (IOException e) { throw new ThingsboardException(e, ThingsboardErrorCode.GENERAL); } if (resource.getResourceType().equals(ResourceType.LWM2M_MODEL) && toLwM2mObject(resource, true) == null) { @@ -194,8 +197,7 @@ public class DefaultTbResourceService implements TbResourceService { if (isSave) { LwM2mResourceObserve lwM2MResourceObserve = new LwM2mResourceObserve(k, v.name, false, false, false); resources.add(lwM2MResourceObserve); - } - else if (v.operations.isReadable()) { + } else if (v.operations.isReadable()) { LwM2mResourceObserve lwM2MResourceObserve = new LwM2mResourceObserve(k, v.name, false, false, false); resources.add(lwM2MResourceObserve); } @@ -204,8 +206,7 @@ public class DefaultTbResourceService implements TbResourceService { instance.setResources(resources.toArray(LwM2mResourceObserve[]::new)); lwM2mObject.setInstances(new LwM2mInstance[]{instance}); return lwM2mObject; - } - else { + } else { return null; } } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java index 05799fc643..d0ed543b76 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/AbstractCleanUpService.java @@ -15,8 +15,13 @@ */ package org.thingsboard.server.service.ttl; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.msg.queue.ServiceType; +import org.thingsboard.server.queue.discovery.PartitionService; import java.sql.Connection; import java.sql.DriverManager; @@ -27,43 +32,12 @@ import java.sql.Statement; @Slf4j +@RequiredArgsConstructor public abstract class AbstractCleanUpService { - @Value("${spring.datasource.url}") - protected String dbUrl; + private final PartitionService partitionService; - @Value("${spring.datasource.username}") - protected String dbUserName; - - @Value("${spring.datasource.password}") - protected String dbPassword; - - protected long executeQuery(Connection conn, String query) throws SQLException { - try (Statement statement = conn.createStatement(); ResultSet resultSet = statement.executeQuery(query)) { - if (log.isDebugEnabled()) { - getWarnings(statement); - } - resultSet.next(); - return resultSet.getLong(1); - } - } - - protected void getWarnings(Statement statement) throws SQLException { - SQLWarning warnings = statement.getWarnings(); - if (warnings != null) { - log.debug("{}", warnings.getMessage()); - SQLWarning nextWarning = warnings.getNextWarning(); - while (nextWarning != null) { - log.debug("{}", nextWarning.getMessage()); - nextWarning = nextWarning.getNextWarning(); - } - } + protected boolean isSystemTenantPartitionMine(){ + return partitionService.resolve(ServiceType.TB_CORE, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID).isMyPartition(); } - - protected abstract void doCleanUp(Connection connection) throws SQLException; - - protected Connection getConnection() throws SQLException { - return DriverManager.getConnection(dbUrl, dbUserName, dbPassword); - } - } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java similarity index 99% rename from application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java rename to application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java index 3b76a6cbca..051a6c92b2 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/alarms/AlarmsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/AlarmsCleanUpService.java @@ -13,7 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.ttl.alarms; +package org.thingsboard.server.service.ttl; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -52,6 +52,7 @@ import java.util.concurrent.TimeUnit; @Slf4j @RequiredArgsConstructor public class AlarmsCleanUpService { + @Value("${sql.ttl.alarms.removal_batch_size}") private Integer removalBatchSize; diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/edge/EdgeEventsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/EdgeEventsCleanUpService.java similarity index 64% rename from application/src/main/java/org/thingsboard/server/service/ttl/edge/EdgeEventsCleanUpService.java rename to application/src/main/java/org/thingsboard/server/service/ttl/EdgeEventsCleanUpService.java index e93a82c7eb..3cdd1f71bd 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/edge/EdgeEventsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/EdgeEventsCleanUpService.java @@ -13,20 +13,18 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.ttl.edge; +package org.thingsboard.server.service.ttl; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; -import org.thingsboard.server.dao.util.PsqlDao; +import org.thingsboard.server.dao.edge.EdgeService; +import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.ttl.AbstractCleanUpService; -import java.sql.Connection; -import java.sql.DriverManager; -import java.sql.SQLException; - -@PsqlDao +@TbCoreComponent @Slf4j @Service public class EdgeEventsCleanUpService extends AbstractCleanUpService { @@ -37,20 +35,18 @@ public class EdgeEventsCleanUpService extends AbstractCleanUpService { @Value("${sql.ttl.edge_events.enabled}") private boolean ttlTaskExecutionEnabled; + private final EdgeService edgeService; + + public EdgeEventsCleanUpService(PartitionService partitionService, EdgeService edgeService) { + super(partitionService); + this.edgeService = edgeService; + } + @Scheduled(initialDelayString = "${sql.ttl.edge_events.execution_interval_ms}", fixedDelayString = "${sql.ttl.edge_events.execution_interval_ms}") public void cleanUp() { - if (ttlTaskExecutionEnabled) { - try (Connection conn = getConnection()) { - doCleanUp(conn); - } catch (SQLException e) { - log.error("SQLException occurred during TTL task execution ", e); - } + if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) { + edgeService.cleanupEvents(ttl); } } - @Override - protected void doCleanUp(Connection connection) throws SQLException { - long totalEdgeEventsRemoved = executeQuery(connection, "call cleanup_edge_events_by_ttl(" + ttl + ", 0);"); - log.info("Total edge events removed by TTL: [{}]", totalEdgeEventsRemoved); - } } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/events/EventsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java similarity index 65% rename from application/src/main/java/org/thingsboard/server/service/ttl/events/EventsCleanUpService.java rename to application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java index 407c88261f..a51910b7ed 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/events/EventsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/EventsCleanUpService.java @@ -13,20 +13,18 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.ttl.events; +package org.thingsboard.server.service.ttl; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; -import org.thingsboard.server.dao.util.PsqlDao; +import org.thingsboard.server.dao.event.EventService; +import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.ttl.AbstractCleanUpService; -import java.sql.Connection; -import java.sql.DriverManager; -import java.sql.SQLException; - -@PsqlDao +@TbCoreComponent @Slf4j @Service public class EventsCleanUpService extends AbstractCleanUpService { @@ -40,20 +38,18 @@ public class EventsCleanUpService extends AbstractCleanUpService { @Value("${sql.ttl.events.enabled}") private boolean ttlTaskExecutionEnabled; + private final EventService eventService; + + public EventsCleanUpService(PartitionService partitionService, EventService eventService) { + super(partitionService); + this.eventService = eventService; + } + @Scheduled(initialDelayString = "${sql.ttl.events.execution_interval_ms}", fixedDelayString = "${sql.ttl.events.execution_interval_ms}") public void cleanUp() { - if (ttlTaskExecutionEnabled) { - try (Connection conn = getConnection()) { - doCleanUp(conn); - } catch (SQLException e) { - log.error("SQLException occurred during TTL task execution ", e); - } + if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) { + eventService.cleanupEvents(ttl, debugTtl); } } - @Override - protected void doCleanUp(Connection connection) throws SQLException { - long totalEventsRemoved = executeQuery(connection, "call cleanup_events_by_ttl(" + ttl + ", " + debugTtl + ", 0);"); - log.info("Total events removed by TTL: [{}]", totalEventsRemoved); - } } \ No newline at end of file diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/AbstractTimeseriesCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/TimeseriesCleanUpService.java similarity index 61% rename from application/src/main/java/org/thingsboard/server/service/ttl/timeseries/AbstractTimeseriesCleanUpService.java rename to application/src/main/java/org/thingsboard/server/service/ttl/TimeseriesCleanUpService.java index ee2d437a22..55c746b580 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/AbstractTimeseriesCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/TimeseriesCleanUpService.java @@ -13,19 +13,21 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.service.ttl.timeseries; +package org.thingsboard.server.service.ttl; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Service; +import org.thingsboard.server.dao.timeseries.TimeseriesService; +import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.ttl.AbstractCleanUpService; -import java.sql.Connection; -import java.sql.DriverManager; -import java.sql.SQLException; - +@TbCoreComponent @Slf4j -public abstract class AbstractTimeseriesCleanUpService extends AbstractCleanUpService { +@Service +public class TimeseriesCleanUpService extends AbstractCleanUpService { @Value("${sql.ttl.ts.ts_key_value_ttl}") protected long systemTtl; @@ -33,14 +35,17 @@ public abstract class AbstractTimeseriesCleanUpService extends AbstractCleanUpSe @Value("${sql.ttl.ts.enabled}") private boolean ttlTaskExecutionEnabled; + private final TimeseriesService timeseriesService; + + public TimeseriesCleanUpService(PartitionService partitionService, TimeseriesService timeseriesService) { + super(partitionService); + this.timeseriesService = timeseriesService; + } + @Scheduled(initialDelayString = "${sql.ttl.ts.execution_interval_ms}", fixedDelayString = "${sql.ttl.ts.execution_interval_ms}") public void cleanUp() { - if (ttlTaskExecutionEnabled) { - try (Connection conn = getConnection()) { - doCleanUp(conn); - } catch (SQLException e) { - log.error("SQLException occurred during TTL task execution ", e); - } + if (ttlTaskExecutionEnabled && isSystemTenantPartitionMine()) { + timeseriesService.cleanup(systemTtl); } } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/PsqlTimeseriesCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/PsqlTimeseriesCleanUpService.java deleted file mode 100644 index 3197f0cbb0..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/PsqlTimeseriesCleanUpService.java +++ /dev/null @@ -1,44 +0,0 @@ -/** - * Copyright © 2016-2021 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.service.ttl.timeseries; - -import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.stereotype.Service; -import org.thingsboard.server.dao.model.ModelConstants; -import org.thingsboard.server.dao.util.PsqlDao; -import org.thingsboard.server.dao.util.SqlTsDao; - -import java.sql.Connection; -import java.sql.SQLException; - -@SqlTsDao -@PsqlDao -@Service -@Slf4j -public class PsqlTimeseriesCleanUpService extends AbstractTimeseriesCleanUpService { - - @Value("${sql.postgres.ts_key_value_partitioning}") - private String partitionType; - - @Override - protected void doCleanUp(Connection connection) throws SQLException { - long totalPartitionsRemoved = executeQuery(connection, "call drop_partitions_by_max_ttl('" + partitionType + "'," + systemTtl + ", 0);"); - log.info("Total partitions removed by TTL: [{}]", totalPartitionsRemoved); - long totalEntitiesTelemetryRemoved = executeQuery(connection, "call cleanup_timeseries_by_ttl('" + ModelConstants.NULL_UUID + "'," + systemTtl + ", 0);"); - log.info("Total telemetry removed stats by TTL for entities: [{}]", totalEntitiesTelemetryRemoved); - } -} \ No newline at end of file diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/TimescaleTimeseriesCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/TimescaleTimeseriesCleanUpService.java deleted file mode 100644 index 0ed61ef97c..0000000000 --- a/application/src/main/java/org/thingsboard/server/service/ttl/timeseries/TimescaleTimeseriesCleanUpService.java +++ /dev/null @@ -1,36 +0,0 @@ -/** - * Copyright © 2016-2021 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.service.ttl.timeseries; - -import lombok.extern.slf4j.Slf4j; -import org.springframework.stereotype.Service; -import org.thingsboard.server.dao.model.ModelConstants; -import org.thingsboard.server.dao.util.TimescaleDBTsDao; - -import java.sql.Connection; -import java.sql.SQLException; - -@TimescaleDBTsDao -@Service -@Slf4j -public class TimescaleTimeseriesCleanUpService extends AbstractTimeseriesCleanUpService { - - @Override - protected void doCleanUp(Connection connection) throws SQLException { - long totalEntitiesTelemetryRemoved = executeQuery(connection, "call cleanup_timeseries_by_ttl('" + ModelConstants.NULL_UUID + "'," + systemTtl + ", 0);"); - log.info("Total telemetry removed stats by TTL for entities: [{}]", totalEntitiesTelemetryRemoved); - } -} diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java index ad156ca3a7..02ac3a56c6 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/AbstractLwM2MIntegrationTest.java @@ -17,24 +17,47 @@ package org.thingsboard.server.transport.lwm2m; import com.fasterxml.jackson.core.type.TypeReference; import org.apache.commons.io.IOUtils; +import org.eclipse.californium.core.network.config.NetworkConfig; +import org.eclipse.leshan.client.object.Security; import org.eclipse.leshan.core.util.Hex; +import org.jetbrains.annotations.NotNull; import org.junit.After; import org.junit.Assert; import org.junit.Before; +import org.springframework.mock.web.MockMultipartFile; +import org.springframework.test.web.servlet.request.MockMultipartHttpServletRequestBuilder; +import org.springframework.test.web.servlet.request.MockMvcRequestBuilders; import org.thingsboard.common.util.JacksonUtil; +import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.data.DeviceProfileProvisionType; import org.thingsboard.server.common.data.DeviceProfileType; import org.thingsboard.server.common.data.DeviceTransportType; +import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.ResourceType; import org.thingsboard.server.common.data.TbResource; +import org.thingsboard.server.common.data.device.credentials.lwm2m.LwM2MClientCredentials; import org.thingsboard.server.common.data.device.profile.DefaultDeviceProfileConfiguration; import org.thingsboard.server.common.data.device.profile.DeviceProfileData; import org.thingsboard.server.common.data.device.profile.DisabledDeviceProfileProvisionConfiguration; import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; +import org.thingsboard.server.common.data.query.EntityData; +import org.thingsboard.server.common.data.query.EntityDataPageLink; +import org.thingsboard.server.common.data.query.EntityDataQuery; +import org.thingsboard.server.common.data.query.EntityKey; +import org.thingsboard.server.common.data.query.EntityKeyType; +import org.thingsboard.server.common.data.query.SingleEntityFilter; +import org.thingsboard.server.common.data.security.DeviceCredentials; +import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.controller.AbstractWebsocketTest; import org.thingsboard.server.controller.TbTestWebSocketClient; import org.thingsboard.server.dao.service.DaoSqlTest; +import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; +import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; +import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; +import org.thingsboard.server.transport.lwm2m.client.LwM2MTestClient; +import org.thingsboard.server.transport.lwm2m.secure.credentials.LwM2MCredentials; import java.io.IOException; import java.io.InputStream; @@ -54,9 +77,15 @@ import java.security.spec.ECPrivateKeySpec; import java.security.spec.ECPublicKeySpec; import java.security.spec.KeySpec; import java.util.Base64; +import java.util.Collections; +import java.util.List; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; +import static org.eclipse.leshan.client.object.Security.noSec; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; +import static org.thingsboard.server.common.data.ota.OtaPackageType.FIRMWARE; + @DaoSqlTest public class AbstractLwM2MIntegrationTest extends AbstractWebsocketTest { @@ -139,6 +168,15 @@ public class AbstractLwM2MIntegrationTest extends AbstractWebsocketTest { // certificates trustedby the server (should contain rootCA) protected final Certificate[] trustedCertificates = new Certificate[1]; + protected static final int SECURE_PORT = 5686; + protected static final NetworkConfig SECURE_COAP_CONFIG = new NetworkConfig().setString("COAP_SECURE_PORT", Integer.toString(SECURE_PORT)); + protected static final String ENDPOINT = "deviceAEndpoint"; + protected static final String SECURE_URI = "coaps://localhost:" + SECURE_PORT; + + protected static final int PORT = 5685; + protected static final Security SECURITY = noSec("coap://localhost:" + PORT, 123); + protected static final NetworkConfig COAP_CONFIG = new NetworkConfig().setString("COAP_PORT", Integer.toString(PORT)); + public AbstractLwM2MIntegrationTest() { // create client credentials try { @@ -262,10 +300,95 @@ public class AbstractLwM2MIntegrationTest extends AbstractWebsocketTest { Assert.assertNotNull(deviceProfile); } + @NotNull + protected Device createDevice(LwM2MClientCredentials clientCredentials) throws Exception { + Device device = new Device(); + device.setName("Device A"); + device.setDeviceProfileId(deviceProfile.getId()); + device.setTenantId(tenantId); + device = doPost("/api/device", device, Device.class); + Assert.assertNotNull(device); + + DeviceCredentials deviceCredentials = + doGet("/api/device/" + device.getId().getId().toString() + "/credentials", DeviceCredentials.class); + Assert.assertEquals(device.getId(), deviceCredentials.getDeviceId()); + deviceCredentials.setCredentialsType(DeviceCredentialsType.LWM2M_CREDENTIALS); + + LwM2MCredentials credentials = new LwM2MCredentials(); + + credentials.setClient(clientCredentials); + + deviceCredentials.setCredentialsValue(JacksonUtil.toString(credentials)); + doPost("/api/device/credentials", deviceCredentials).andExpect(status().isOk()); + return device; + } + + + protected OtaPackageInfo createFirmware() throws Exception { + String CHECKSUM = "4bf5122f344554c53bde2ebb8cd2b7e3d1600ad631c385a5d7cce23c7785459a"; + + OtaPackageInfo firmwareInfo = new OtaPackageInfo(); + firmwareInfo.setDeviceProfileId(deviceProfile.getId()); + firmwareInfo.setType(FIRMWARE); + firmwareInfo.setTitle("My firmware"); + firmwareInfo.setVersion("v1.0"); + + OtaPackageInfo savedFirmwareInfo = doPost("/api/otaPackage", firmwareInfo, OtaPackageInfo.class); + + MockMultipartFile testData = new MockMultipartFile("file", "filename.txt", "text/plain", new byte[]{1}); + + return savaData("/api/otaPackage/" + savedFirmwareInfo.getId().getId().toString() + "?checksum={checksum}&checksumAlgorithm={checksumAlgorithm}", testData, CHECKSUM, "SHA256"); + } + + protected OtaPackageInfo savaData(String urlTemplate, MockMultipartFile content, String... params) throws Exception { + MockMultipartHttpServletRequestBuilder postRequest = MockMvcRequestBuilders.multipart(urlTemplate, params); + postRequest.file(content); + setJwtToken(postRequest); + return readResponse(mockMvc.perform(postRequest).andExpect(status().isOk()), OtaPackageInfo.class); + } + @After public void after() { executor.shutdownNow(); wsClient.close(); } + public void basicTestConnectionObserveTelemetry(Security security, + LwM2MClientCredentials credentials, + NetworkConfig coapConfig, + String endpoint) throws Exception { + createDeviceProfile(TRANSPORT_CONFIGURATION); + Device device = createDevice(credentials); + + SingleEntityFilter sef = new SingleEntityFilter(); + sef.setSingleEntity(device.getId()); + LatestValueCmd latestCmd = new LatestValueCmd(); + latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, "batteryLevel"))); + EntityDataQuery edq = new EntityDataQuery(sef, new EntityDataPageLink(1, 0, null, null), + Collections.emptyList(), Collections.emptyList(), Collections.emptyList()); + + EntityDataCmd cmd = new EntityDataCmd(1, edq, null, latestCmd, null); + TelemetryPluginCmdsWrapper wrapper = new TelemetryPluginCmdsWrapper(); + wrapper.setEntityDataCmds(Collections.singletonList(cmd)); + + wsClient.send(mapper.writeValueAsString(wrapper)); + wsClient.waitForReply(); + + wsClient.registerWaitForUpdate(); + LwM2MTestClient client = new LwM2MTestClient(executor, endpoint); + + client.init(security, coapConfig); + String msg = wsClient.waitForUpdate(); + + EntityDataUpdate update = mapper.readValue(msg, EntityDataUpdate.class); + Assert.assertEquals(1, update.getCmdId()); + List eData = update.getUpdate(); + Assert.assertNotNull(eData); + Assert.assertEquals(1, eData.size()); + Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); + Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES)); + var tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("batteryLevel"); + Assert.assertEquals(42, Long.parseLong(tsValue.getValue())); + client.destroy(); + } } diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java index d98a775c8c..8c20eace51 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/NoSecLwM2MIntegrationTest.java @@ -15,14 +15,8 @@ */ package org.thingsboard.server.transport.lwm2m; -import org.eclipse.californium.core.network.config.NetworkConfig; -import org.eclipse.leshan.client.object.Security; import org.junit.Assert; import org.junit.Test; -import org.springframework.mock.web.MockMultipartFile; -import org.springframework.test.web.servlet.request.MockMultipartHttpServletRequestBuilder; -import org.springframework.test.web.servlet.request.MockMvcRequestBuilders; -import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.OtaPackageInfo; import org.thingsboard.server.common.data.device.credentials.lwm2m.NoSecClientCredentials; @@ -32,116 +26,30 @@ import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.data.query.EntityKey; import org.thingsboard.server.common.data.query.EntityKeyType; import org.thingsboard.server.common.data.query.SingleEntityFilter; -import org.thingsboard.server.common.data.security.DeviceCredentials; -import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; import org.thingsboard.server.transport.lwm2m.client.LwM2MTestClient; -import org.thingsboard.server.transport.lwm2m.secure.credentials.LwM2MCredentials; import java.util.Collections; import java.util.List; -import static org.eclipse.leshan.client.object.Security.noSec; -import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; -import static org.thingsboard.server.common.data.ota.OtaPackageType.FIRMWARE; - public class NoSecLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest { - private final int PORT = 5685; - private final Security SECURITY = noSec("coap://localhost:" + PORT, 123); - private final NetworkConfig COAP_CONFIG = new NetworkConfig().setString("COAP_PORT", Integer.toString(PORT)); - private final String ENDPOINT = "noSecEndpoint"; - - private Device createDevice() throws Exception { - Device device = new Device(); - device.setName("Device A"); - device.setDeviceProfileId(deviceProfile.getId()); - device.setTenantId(tenantId); - device = doPost("/api/device", device, Device.class); - Assert.assertNotNull(device); - - DeviceCredentials deviceCredentials = - doGet("/api/device/" + device.getId().getId().toString() + "/credentials", DeviceCredentials.class); - Assert.assertEquals(device.getId(), deviceCredentials.getDeviceId()); - deviceCredentials.setCredentialsType(DeviceCredentialsType.LWM2M_CREDENTIALS); - - LwM2MCredentials noSecCredentials = new LwM2MCredentials(); - NoSecClientCredentials clientCredentials = new NoSecClientCredentials(); - clientCredentials.setEndpoint(ENDPOINT); - noSecCredentials.setClient(clientCredentials); - deviceCredentials.setCredentialsValue(JacksonUtil.toString(noSecCredentials)); - doPost("/api/device/credentials", deviceCredentials).andExpect(status().isOk()); - return device; - } - - private OtaPackageInfo createFirmware() throws Exception { - String CHECKSUM = "4bf5122f344554c53bde2ebb8cd2b7e3d1600ad631c385a5d7cce23c7785459a"; - - OtaPackageInfo firmwareInfo = new OtaPackageInfo(); - firmwareInfo.setDeviceProfileId(deviceProfile.getId()); - firmwareInfo.setType(FIRMWARE); - firmwareInfo.setTitle("My firmware"); - firmwareInfo.setVersion("v1.0"); - - OtaPackageInfo savedFirmwareInfo = doPost("/api/otaPackage", firmwareInfo, OtaPackageInfo.class); - - MockMultipartFile testData = new MockMultipartFile("file", "filename.txt", "text/plain", new byte[]{1}); - - return savaData("/api/otaPackage/" + savedFirmwareInfo.getId().getId().toString() + "?checksum={checksum}&checksumAlgorithm={checksumAlgorithm}", testData, CHECKSUM, "SHA256"); - } - - protected OtaPackageInfo savaData(String urlTemplate, MockMultipartFile content, String... params) throws Exception { - MockMultipartHttpServletRequestBuilder postRequest = MockMvcRequestBuilders.multipart(urlTemplate, params); - postRequest.file(content); - setJwtToken(postRequest); - return readResponse(mockMvc.perform(postRequest).andExpect(status().isOk()), OtaPackageInfo.class); - } - @Test public void testConnectAndObserveTelemetry() throws Exception { - createDeviceProfile(TRANSPORT_CONFIGURATION); - - Device device = createDevice(); - - SingleEntityFilter sef = new SingleEntityFilter(); - sef.setSingleEntity(device.getId()); - LatestValueCmd latestCmd = new LatestValueCmd(); - latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, "batteryLevel"))); - EntityDataQuery edq = new EntityDataQuery(sef, new EntityDataPageLink(1, 0, null, null), - Collections.emptyList(), Collections.emptyList(), Collections.emptyList()); - - EntityDataCmd cmd = new EntityDataCmd(1, edq, null, latestCmd, null); - TelemetryPluginCmdsWrapper wrapper = new TelemetryPluginCmdsWrapper(); - wrapper.setEntityDataCmds(Collections.singletonList(cmd)); - - wsClient.send(mapper.writeValueAsString(wrapper)); - wsClient.waitForReply(); - - wsClient.registerWaitForUpdate(); - LwM2MTestClient client = new LwM2MTestClient(executor, ENDPOINT); - client.init(SECURITY, COAP_CONFIG); - String msg = wsClient.waitForUpdate(); - - EntityDataUpdate update = mapper.readValue(msg, EntityDataUpdate.class); - Assert.assertEquals(1, update.getCmdId()); - List eData = update.getUpdate(); - Assert.assertNotNull(eData); - Assert.assertEquals(1, eData.size()); - Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); - Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES)); - var tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("batteryLevel"); - Assert.assertEquals(42, Long.parseLong(tsValue.getValue())); - client.destroy(); + NoSecClientCredentials clientCredentials = new NoSecClientCredentials(); + clientCredentials.setEndpoint(ENDPOINT); + super.basicTestConnectionObserveTelemetry(SECURITY, clientCredentials, COAP_CONFIG, ENDPOINT); } @Test public void testFirmwareUpdateWithClientWithoutFirmwareInfo() throws Exception { createDeviceProfile(TRANSPORT_CONFIGURATION); - - Device device = createDevice(); + NoSecClientCredentials clientCredentials = new NoSecClientCredentials(); + clientCredentials.setEndpoint(ENDPOINT); + Device device = createDevice(clientCredentials); OtaPackageInfo firmware = createFirmware(); diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/PskLwm2mIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/PskLwm2mIntegrationTest.java new file mode 100644 index 0000000000..6d3e544d57 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/PskLwm2mIntegrationTest.java @@ -0,0 +1,43 @@ +/** + * Copyright © 2016-2021 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.transport.lwm2m; + +import org.eclipse.leshan.client.object.Security; +import org.eclipse.leshan.core.util.Hex; +import org.junit.Test; +import org.thingsboard.server.common.data.device.credentials.lwm2m.PSKClientCredentials; + +import java.nio.charset.StandardCharsets; + +import static org.eclipse.leshan.client.object.Security.psk; + +public class PskLwm2mIntegrationTest extends AbstractLwM2MIntegrationTest { + + @Test + public void testConnectWithPSKAndObserveTelemetry() throws Exception { + String pskIdentity = "SOME_PSK_ID"; + String pskKey = "73656372657450534b"; + PSKClientCredentials clientCredentials = new PSKClientCredentials(); + clientCredentials.setEndpoint(ENDPOINT); + clientCredentials.setKey(pskKey); + clientCredentials.setIdentity(pskIdentity); + Security security = psk(SECURE_URI, + 123, + pskIdentity.getBytes(StandardCharsets.UTF_8), + Hex.decodeHex(pskKey.toCharArray())); + super.basicTestConnectionObserveTelemetry(security, clientCredentials, SECURE_COAP_CONFIG, ENDPOINT); + } +} diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/RpkLwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/RpkLwM2MIntegrationTest.java new file mode 100644 index 0000000000..eb81f91e2c --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/RpkLwM2MIntegrationTest.java @@ -0,0 +1,39 @@ +/** + * Copyright © 2016-2021 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.transport.lwm2m; + +import org.eclipse.leshan.client.object.Security; +import org.eclipse.leshan.core.util.Hex; +import org.junit.Test; +import org.thingsboard.server.common.data.device.credentials.lwm2m.RPKClientCredentials; + +import static org.eclipse.leshan.client.object.Security.rpk; + +public class RpkLwM2MIntegrationTest extends AbstractLwM2MIntegrationTest { + @Test + public void testConnectWithRPKAndObserveTelemetry() throws Exception { + RPKClientCredentials rpkClientCredentials = new RPKClientCredentials(); + rpkClientCredentials.setEndpoint(ENDPOINT); + rpkClientCredentials.setKey(Hex.encodeHexString(clientPublicKey.getEncoded())); + Security security = rpk(SECURE_URI, + 123, + clientPublicKey.getEncoded(), + clientPrivateKey.getEncoded(), + serverX509Cert.getPublicKey().getEncoded()); + super.basicTestConnectionObserveTelemetry(security, rpkClientCredentials, SECURE_COAP_CONFIG, ENDPOINT); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/transport/lwm2m/X509LwM2MIntegrationTest.java b/application/src/test/java/org/thingsboard/server/transport/lwm2m/X509LwM2MIntegrationTest.java index 661f7c5474..1760c16f02 100644 --- a/application/src/test/java/org/thingsboard/server/transport/lwm2m/X509LwM2MIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/transport/lwm2m/X509LwM2MIntegrationTest.java @@ -15,147 +15,38 @@ */ package org.thingsboard.server.transport.lwm2m; -import org.eclipse.californium.core.network.config.NetworkConfig; import org.eclipse.leshan.client.object.Security; -import org.jetbrains.annotations.NotNull; -import org.junit.Assert; -import org.junit.Ignore; import org.junit.Test; -import org.thingsboard.common.util.JacksonUtil; -import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.device.credentials.lwm2m.X509ClientCredentials; -import org.thingsboard.server.common.data.query.EntityData; -import org.thingsboard.server.common.data.query.EntityDataPageLink; -import org.thingsboard.server.common.data.query.EntityDataQuery; -import org.thingsboard.server.common.data.query.EntityKey; -import org.thingsboard.server.common.data.query.EntityKeyType; -import org.thingsboard.server.common.data.query.SingleEntityFilter; -import org.thingsboard.server.common.data.security.DeviceCredentials; -import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.transport.util.SslUtil; -import org.thingsboard.server.service.telemetry.cmd.TelemetryPluginCmdsWrapper; -import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; -import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; -import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; -import org.thingsboard.server.transport.lwm2m.client.LwM2MTestClient; -import org.thingsboard.server.transport.lwm2m.secure.credentials.LwM2MCredentials; - -import java.util.Collections; -import java.util.List; import static org.eclipse.leshan.client.object.Security.x509; -import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; public class X509LwM2MIntegrationTest extends AbstractLwM2MIntegrationTest { - private final int port = 5686; - private final NetworkConfig coapConfig = new NetworkConfig().setString("COAP_SECURE_PORT", Integer.toString(port)); - private final String endpoint = "deviceAEndpoint"; - private final String serverUri = "coaps://localhost:" + port; - - private Device createDevice(X509ClientCredentials clientCredentials) throws Exception { - Device device = new Device(); - device.setName("Device A"); - device.setDeviceProfileId(deviceProfile.getId()); - device.setTenantId(tenantId); - device = doPost("/api/device", device, Device.class); - Assert.assertNotNull(device); - - DeviceCredentials deviceCredentials = - doGet("/api/device/" + device.getId().getId().toString() + "/credentials", DeviceCredentials.class); - Assert.assertEquals(device.getId(), deviceCredentials.getDeviceId()); - deviceCredentials.setCredentialsType(DeviceCredentialsType.LWM2M_CREDENTIALS); - - LwM2MCredentials credentials = new LwM2MCredentials(); - - credentials.setClient(clientCredentials); - - deviceCredentials.setCredentialsValue(JacksonUtil.toString(credentials)); - doPost("/api/device/credentials", deviceCredentials).andExpect(status().isOk()); - return device; - } - - //TODO: use different endpoints to isolate tests. - @Ignore() @Test public void testConnectAndObserveTelemetry() throws Exception { - createDeviceProfile(TRANSPORT_CONFIGURATION); X509ClientCredentials credentials = new X509ClientCredentials(); - credentials.setEndpoint(endpoint+1); - Device device = createDevice(credentials); - - SingleEntityFilter sef = new SingleEntityFilter(); - sef.setSingleEntity(device.getId()); - LatestValueCmd latestCmd = new LatestValueCmd(); - latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, "batteryLevel"))); - EntityDataQuery edq = new EntityDataQuery(sef, new EntityDataPageLink(1, 0, null, null), - Collections.emptyList(), Collections.emptyList(), Collections.emptyList()); - - EntityDataCmd cmd = new EntityDataCmd(1, edq, null, latestCmd, null); - TelemetryPluginCmdsWrapper wrapper = new TelemetryPluginCmdsWrapper(); - wrapper.setEntityDataCmds(Collections.singletonList(cmd)); - - wsClient.send(mapper.writeValueAsString(wrapper)); - wsClient.waitForReply(); - - wsClient.registerWaitForUpdate(); - LwM2MTestClient client = new LwM2MTestClient(executor, endpoint+1); - Security security = x509(serverUri, 123, clientX509Cert.getEncoded(), clientPrivateKeyFromCert.getEncoded(), serverX509Cert.getEncoded()); - client.init(security, coapConfig); - String msg = wsClient.waitForUpdate(); - - EntityDataUpdate update = mapper.readValue(msg, EntityDataUpdate.class); - Assert.assertEquals(1, update.getCmdId()); - List eData = update.getUpdate(); - Assert.assertNotNull(eData); - Assert.assertEquals(1, eData.size()); - Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); - Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES)); - var tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("batteryLevel"); - Assert.assertEquals(42, Long.parseLong(tsValue.getValue())); - client.destroy(); + credentials.setEndpoint(ENDPOINT); + Security security = x509(SECURE_URI, + 123, + clientX509Cert.getEncoded(), + clientPrivateKeyFromCert.getEncoded(), + serverX509Cert.getEncoded()); + super.basicTestConnectionObserveTelemetry(security, credentials, SECURE_COAP_CONFIG, ENDPOINT); } @Test public void testConnectWithCertAndObserveTelemetry() throws Exception { - createDeviceProfile(TRANSPORT_CONFIGURATION); X509ClientCredentials credentials = new X509ClientCredentials(); - credentials.setEndpoint(endpoint); + credentials.setEndpoint(ENDPOINT); credentials.setCert(SslUtil.getCertificateString(clientX509CertNotTrusted)); - Device device = createDevice(credentials); - - SingleEntityFilter sef = new SingleEntityFilter(); - sef.setSingleEntity(device.getId()); - LatestValueCmd latestCmd = new LatestValueCmd(); - latestCmd.setKeys(Collections.singletonList(new EntityKey(EntityKeyType.TIME_SERIES, "batteryLevel"))); - EntityDataQuery edq = new EntityDataQuery(sef, new EntityDataPageLink(1, 0, null, null), - Collections.emptyList(), Collections.emptyList(), Collections.emptyList()); - - EntityDataCmd cmd = new EntityDataCmd(1, edq, null, latestCmd, null); - TelemetryPluginCmdsWrapper wrapper = new TelemetryPluginCmdsWrapper(); - wrapper.setEntityDataCmds(Collections.singletonList(cmd)); - - wsClient.send(mapper.writeValueAsString(wrapper)); - wsClient.waitForReply(); - - wsClient.registerWaitForUpdate(); - LwM2MTestClient client = new LwM2MTestClient(executor, endpoint); - - Security security = x509(serverUri, 123, clientX509CertNotTrusted.getEncoded(), clientPrivateKeyFromCert.getEncoded(), serverX509Cert.getEncoded()); - - client.init(security, coapConfig); - String msg = wsClient.waitForUpdate(); - - EntityDataUpdate update = mapper.readValue(msg, EntityDataUpdate.class); - Assert.assertEquals(1, update.getCmdId()); - List eData = update.getUpdate(); - Assert.assertNotNull(eData); - Assert.assertEquals(1, eData.size()); - Assert.assertEquals(device.getId(), eData.get(0).getEntityId()); - Assert.assertNotNull(eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES)); - var tsValue = eData.get(0).getLatest().get(EntityKeyType.TIME_SERIES).get("batteryLevel"); - Assert.assertEquals(42, Long.parseLong(tsValue.getValue())); - client.destroy(); + Security security = x509(SECURE_URI, + 123, + clientX509CertNotTrusted.getEncoded(), + clientPrivateKeyFromCert.getEncoded(), + serverX509Cert.getEncoded()); + super.basicTestConnectionObserveTelemetry(security, credentials, SECURE_COAP_CONFIG, ENDPOINT); } } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java index a7b9145c01..53472fd261 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/edge/EdgeService.java @@ -93,4 +93,6 @@ public interface EdgeService { Object activateInstance(String licenseSecret, String releaseDate); String findMissingToRelatedRuleChains(TenantId tenantId, EdgeId edgeId); + + void cleanupEvents(long ttl); } 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 ea25568375..db1c77697e 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 @@ -46,4 +46,6 @@ public interface EventService { void removeEvents(TenantId tenantId, EntityId entityId); + void cleanupEvents(long ttl, long debugTtl); + } diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java index 6fcb5ca2ca..b1a2541fc7 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java @@ -52,4 +52,6 @@ public interface TimeseriesService { List findAllKeysByDeviceProfileId(TenantId tenantId, DeviceProfileId deviceProfileId); List findAllKeysByEntityIds(TenantId tenantId, List entityIds); + + void cleanup(long systemTtl); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/HasKey.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/AbstractLwM2MClientCredentialsWithKey.java similarity index 63% rename from common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/HasKey.java rename to common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/AbstractLwM2MClientCredentialsWithKey.java index ec62765298..2a3c0ab434 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/HasKey.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/AbstractLwM2MClientCredentialsWithKey.java @@ -15,20 +15,25 @@ */ package org.thingsboard.server.common.data.device.credentials.lwm2m; +import com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.Getter; +import lombok.Setter; import lombok.SneakyThrows; import org.apache.commons.codec.binary.Hex; -public abstract class HasKey extends AbstractLwM2MClientCredentials { - private byte[] key; +public abstract class AbstractLwM2MClientCredentialsWithKey extends AbstractLwM2MClientCredentials { + @Getter + @Setter + private String key; + + private byte[] keyInBytes; @SneakyThrows - public void setKey(String key) { - if (key != null) { - this.key = Hex.decodeHex(key.toLowerCase().toCharArray()); + @JsonIgnore + public byte[] getDecodedKey() { + if (keyInBytes == null) { + keyInBytes = Hex.decodeHex(key.toLowerCase().toCharArray()); } - } - - public byte[] getKey() { - return key; + return keyInBytes; } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/PSKClientCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/PSKClientCredentials.java index 2566af7da8..f90e85ff4a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/PSKClientCredentials.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/PSKClientCredentials.java @@ -20,7 +20,7 @@ import lombok.Setter; @Getter @Setter -public class PSKClientCredentials extends HasKey { +public class PSKClientCredentials extends AbstractLwM2MClientCredentialsWithKey { private String identity; @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/RPKClientCredentials.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/RPKClientCredentials.java index fe329558f8..4ebe2de71a 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/RPKClientCredentials.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/credentials/lwm2m/RPKClientCredentials.java @@ -15,7 +15,7 @@ */ package org.thingsboard.server.common.data.device.credentials.lwm2m; -public class RPKClientCredentials extends HasKey { +public class RPKClientCredentials extends AbstractLwM2MClientCredentialsWithKey { @Override public LwM2MSecurityMode getSecurityConfigClientMode() { diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java index 8fff78ae86..56d49fc779 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java @@ -64,6 +64,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; @Slf4j public class CoapTransportResource extends AbstractCoapTransportResource { @@ -75,9 +76,8 @@ public class CoapTransportResource extends AbstractCoapTransportResource { private static final int REQUEST_ID_POSITION_CERTIFICATE_REQUEST = 4; private static final String DTLS_SESSION_ID_KEY = "DTLS_SESSION_ID"; - private final ConcurrentMap tokenToSessionInfoMap = new ConcurrentHashMap<>(); - private final ConcurrentMap tokenToObserveNotificationSeqMap = new ConcurrentHashMap<>(); - private final ConcurrentMap sessionInfoToObserveRelationMap = new ConcurrentHashMap<>(); + private final ConcurrentMap tokenToCoapSessionInfoMap = new ConcurrentHashMap<>(); + private final ConcurrentMap sessionInfoToObserveRelationMap = new ConcurrentHashMap<>(); private final Set rpcSubscriptions = ConcurrentHashMap.newKeySet(); private final Set attributeSubscriptions = ConcurrentHashMap.newKeySet(); @@ -93,7 +93,11 @@ public class CoapTransportResource extends AbstractCoapTransportResource { this.timeout = coapServerService.getTimeout(); this.sessionReportTimeout = ctx.getSessionReportTimeout(); ctx.getScheduler().scheduleAtFixedRate(() -> { - Set observeSessions = sessionInfoToObserveRelationMap.keySet(); + Set coapObserveSessionInfos = sessionInfoToObserveRelationMap.keySet(); + Set observeSessions = coapObserveSessionInfos + .stream() + .map(CoapObserveSessionInfo::getSessionInfoProto) + .collect(Collectors.toSet()); observeSessions.forEach(this::reportActivity); }, new Random().nextInt((int) sessionReportTimeout), sessionReportTimeout, TimeUnit.MILLISECONDS); } @@ -111,17 +115,17 @@ public class CoapTransportResource extends AbstractCoapTransportResource { relation.setEstablished(); addObserveRelation(relation); } - AtomicInteger notificationCounter = tokenToObserveNotificationSeqMap.computeIfAbsent(token, s -> new AtomicInteger(0)); - response.getOptions().setObserve(notificationCounter.getAndIncrement()); + AtomicInteger observeNotificationCounter = tokenToCoapSessionInfoMap.get(token).getObserveNotificationCounter(); + response.getOptions().setObserve(observeNotificationCounter.getAndIncrement()); } // ObserveLayer takes care of the else case } - public void clearAndNotifyObserveRelation(ObserveRelation relation, CoAP.ResponseCode code) { + private void clearAndNotifyObserveRelation(ObserveRelation relation, CoAP.ResponseCode code) { relation.cancel(); relation.getExchange().sendResponse(new Response(code)); } - public Map getSessionInfoToObserveRelationMap() { + private Map getCoapSessionInfoToObserveRelationMap() { return sessionInfoToObserveRelationMap; } @@ -277,8 +281,8 @@ public class CoapTransportResource extends AbstractCoapTransportResource { new CoapOkCallback(exchange, CoAP.ResponseCode.CREATED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); break; case SUBSCRIBE_ATTRIBUTES_REQUEST: - TransportProtos.SessionInfoProto currentAttrSession = tokenToSessionInfoMap.get(getTokenFromRequest(request)); - if (currentAttrSession == null) { + CoapObserveSessionInfo currentCoapObserveAttrSessionInfo = tokenToCoapSessionInfoMap.get(getTokenFromRequest(request)); + if (currentCoapObserveAttrSessionInfo == null) { attributeSubscriptions.add(sessionId); registerAsyncCoapSession(exchange, sessionInfo, coapTransportAdaptor, transportConfigurationContainer.getRpcRequestDynamicMessageBuilder(), getTokenFromRequest(request)); @@ -290,20 +294,20 @@ public class CoapTransportResource extends AbstractCoapTransportResource { } break; case UNSUBSCRIBE_ATTRIBUTES_REQUEST: - TransportProtos.SessionInfoProto attrSession = lookupAsyncSessionInfo(getTokenFromRequest(request)); - if (attrSession != null) { + CoapObserveSessionInfo coapObserveAttrSessionInfo = lookupAsyncSessionInfo(getTokenFromRequest(request)); + if (coapObserveAttrSessionInfo != null) { + TransportProtos.SessionInfoProto attrSession = coapObserveAttrSessionInfo.getSessionInfoProto(); UUID attrSessionId = toSessionId(attrSession); attributeSubscriptions.remove(attrSessionId); - sessionInfoToObserveRelationMap.remove(attrSession); transportService.process(attrSession, TransportProtos.SubscribeToAttributeUpdatesMsg.newBuilder().setUnsubscribe(true).build(), new CoapOkCallback(exchange, CoAP.ResponseCode.DELETED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); - closeAndDeregister(sessionInfo); } + closeAndDeregister(sessionInfo); break; case SUBSCRIBE_RPC_COMMANDS_REQUEST: - TransportProtos.SessionInfoProto currentRpcSession = tokenToSessionInfoMap.get(getTokenFromRequest(request)); - if (currentRpcSession == null) { + CoapObserveSessionInfo currentCoapObserveRpcSessionInfo = tokenToCoapSessionInfoMap.get(getTokenFromRequest(request)); + if (currentCoapObserveRpcSessionInfo == null) { rpcSubscriptions.add(sessionId); registerAsyncCoapSession(exchange, sessionInfo, coapTransportAdaptor, transportConfigurationContainer.getRpcRequestDynamicMessageBuilder(), getTokenFromRequest(request)); @@ -314,16 +318,16 @@ public class CoapTransportResource extends AbstractCoapTransportResource { } break; case UNSUBSCRIBE_RPC_COMMANDS_REQUEST: - TransportProtos.SessionInfoProto rpcSession = lookupAsyncSessionInfo(getTokenFromRequest(request)); - if (rpcSession != null) { + CoapObserveSessionInfo coapObserveRpcSessionInfo = lookupAsyncSessionInfo(getTokenFromRequest(request)); + if (coapObserveRpcSessionInfo != null) { + TransportProtos.SessionInfoProto rpcSession = coapObserveRpcSessionInfo.getSessionInfoProto(); UUID rpcSessionId = toSessionId(rpcSession); rpcSubscriptions.remove(rpcSessionId); - sessionInfoToObserveRelationMap.remove(rpcSession); transportService.process(rpcSession, TransportProtos.SubscribeToRPCMsg.newBuilder().setUnsubscribe(true).build(), new CoapOkCallback(exchange, CoAP.ResponseCode.DELETED, CoAP.ResponseCode.INTERNAL_SERVER_ERROR)); - closeAndDeregister(sessionInfo); } + closeAndDeregister(sessionInfo); break; case TO_DEVICE_RPC_RESPONSE: transportService.process(sessionInfo, @@ -355,13 +359,12 @@ public class CoapTransportResource extends AbstractCoapTransportResource { return new UUID(sessionInfoProto.getSessionIdMSB(), sessionInfoProto.getSessionIdLSB()); } - private TransportProtos.SessionInfoProto lookupAsyncSessionInfo(String token) { - tokenToObserveNotificationSeqMap.remove(token); - return tokenToSessionInfoMap.remove(token); + private CoapObserveSessionInfo lookupAsyncSessionInfo(String token) { + return tokenToCoapSessionInfoMap.remove(token); } private void registerAsyncCoapSession(CoapExchange exchange, TransportProtos.SessionInfoProto sessionInfo, CoapTransportAdaptor coapTransportAdaptor, DynamicMessage.Builder rpcRequestDynamicMessageBuilder, String token) { - tokenToSessionInfoMap.putIfAbsent(token, sessionInfo); + tokenToCoapSessionInfoMap.putIfAbsent(token, new CoapObserveSessionInfo(sessionInfo)); transportService.registerAsyncSession(sessionInfo, getCoapSessionListener(exchange, coapTransportAdaptor, rpcRequestDynamicMessageBuilder, sessionInfo)); transportService.process(sessionInfo, getSessionEventMsg(TransportProtos.SessionEvent.OPEN), null); } @@ -476,45 +479,40 @@ public class CoapTransportResource extends AbstractCoapTransportResource { } @Override - public void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg) { + public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg msg) { + log.trace("[{}] Received attributes update notification to device", sessionId); try { exchange.respond(coapTransportAdaptor.convertToPublish(isConRequest(), msg)); } catch (AdaptorException e) { log.trace("Failed to reply due to error", e); - exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR); + closeObserveRelationAndNotify(sessionId, CoAP.ResponseCode.INTERNAL_SERVER_ERROR); + closeAndDeregister(); } } @Override public void onRemoteSessionCloseCommand(UUID sessionId, TransportProtos.SessionCloseNotificationProto sessionCloseNotification) { log.trace("[{}] Received the remote command to close the session: {}", sessionId, sessionCloseNotification.getMessage()); - Map sessionToObserveRelationMap = coapTransportResource.getSessionInfoToObserveRelationMap(); - if (coapTransportResource.getObserverCount() > 0 && !CollectionUtils.isEmpty(sessionToObserveRelationMap)) { - Set observeSessions = sessionToObserveRelationMap.keySet(); - Optional observeSessionToClose = observeSessions.stream().filter(sessionInfoProto -> { - UUID observeSessionId = new UUID(sessionInfoProto.getSessionIdMSB(), sessionInfoProto.getSessionIdLSB()); - return observeSessionId.equals(sessionId); - }).findFirst(); - if (observeSessionToClose.isPresent()) { - TransportProtos.SessionInfoProto sessionInfoProto = observeSessionToClose.get(); - ObserveRelation observeRelation = sessionToObserveRelationMap.get(sessionInfoProto); - coapTransportResource.clearAndNotifyObserveRelation(observeRelation, CoAP.ResponseCode.SERVICE_UNAVAILABLE); - } - } + closeObserveRelationAndNotify(sessionId, CoAP.ResponseCode.SERVICE_UNAVAILABLE); + closeAndDeregister(); } @Override - public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg msg) { - boolean successful; + public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg msg) { + log.trace("[{}] Received RPC command to device", sessionId); + boolean successful = true; try { exchange.respond(coapTransportAdaptor.convertToPublish(isConRequest(), msg, rpcRequestDynamicMessageBuilder)); - successful = true; } catch (AdaptorException e) { log.trace("Failed to reply due to error", e); - exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR); + closeObserveRelationAndNotify(sessionId, CoAP.ResponseCode.INTERNAL_SERVER_ERROR); successful = false; + } finally { + coapTransportResource.transportService.process(sessionInfo, msg, !successful, TransportServiceCallback.EMPTY); + if (!successful) { + closeAndDeregister(); + } } - coapTransportResource.transportService.process(sessionInfo, msg, !successful, TransportServiceCallback.EMPTY); } @Override @@ -530,6 +528,30 @@ public class CoapTransportResource extends AbstractCoapTransportResource { private boolean isConRequest() { return exchange.advanced().getRequest().isConfirmable(); } + + private void closeObserveRelationAndNotify(UUID sessionId, CoAP.ResponseCode responseCode) { + Map sessionToObserveRelationMap = coapTransportResource.getCoapSessionInfoToObserveRelationMap(); + if (coapTransportResource.getObserverCount() > 0 && !CollectionUtils.isEmpty(sessionToObserveRelationMap)) { + Optional observeSessionToClose = sessionToObserveRelationMap.keySet().stream().filter(coapObserveSessionInfo -> { + TransportProtos.SessionInfoProto sessionToDelete = coapObserveSessionInfo.getSessionInfoProto(); + UUID observeSessionId = new UUID(sessionToDelete.getSessionIdMSB(), sessionToDelete.getSessionIdLSB()); + return observeSessionId.equals(sessionId); + }).findFirst(); + if (observeSessionToClose.isPresent()) { + CoapObserveSessionInfo coapObserveSessionInfo = observeSessionToClose.get(); + ObserveRelation observeRelation = sessionToObserveRelationMap.get(coapObserveSessionInfo); + coapTransportResource.clearAndNotifyObserveRelation(observeRelation, responseCode); + } + } + } + + private void closeAndDeregister() { + Request request = exchange.advanced().getRequest(); + String token = coapTransportResource.getTokenFromRequest(request); + CoapObserveSessionInfo deleted = coapTransportResource.lookupAsyncSessionInfo(token); + coapTransportResource.closeAndDeregister(deleted.getSessionInfoProto()); + } + } public class CoapResourceObserver implements ResourceObserver { @@ -554,7 +576,7 @@ public class CoapTransportResource extends AbstractCoapTransportResource { public void addedObserveRelation(ObserveRelation relation) { Request request = relation.getExchange().getRequest(); String token = getTokenFromRequest(request); - sessionInfoToObserveRelationMap.putIfAbsent(tokenToSessionInfoMap.get(token), relation); + sessionInfoToObserveRelationMap.putIfAbsent(tokenToCoapSessionInfoMap.get(token), relation); log.trace("Added Observe relation for token: {}", token); } @@ -562,8 +584,7 @@ public class CoapTransportResource extends AbstractCoapTransportResource { public void removedObserveRelation(ObserveRelation relation) { Request request = relation.getExchange().getRequest(); String token = getTokenFromRequest(request); - TransportProtos.SessionInfoProto session = tokenToSessionInfoMap.get(token); - sessionInfoToObserveRelationMap.remove(session); + sessionInfoToObserveRelationMap.remove(tokenToCoapSessionInfoMap.get(token)); log.trace("Relation removed for token: {}", token); } } @@ -574,7 +595,6 @@ public class CoapTransportResource extends AbstractCoapTransportResource { transportService.deregisterSession(session); rpcSubscriptions.remove(sessionId); attributeSubscriptions.remove(sessionId); - sessionInfoToObserveRelationMap.remove(session); } private TransportConfigurationContainer getTransportConfigurationContainer(DeviceProfile deviceProfile) throws AdaptorException { @@ -640,4 +660,17 @@ public class CoapTransportResource extends AbstractCoapTransportResource { this.jsonPayload = jsonPayload; } } + + @Data + private static class CoapObserveSessionInfo { + + private final TransportProtos.SessionInfoProto sessionInfoProto; + private final AtomicInteger observeNotificationCounter; + + private CoapObserveSessionInfo(TransportProtos.SessionInfoProto sessionInfoProto) { + this.sessionInfoProto = sessionInfoProto; + this.observeNotificationCounter = new AtomicInteger(0); + } + } + } diff --git a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java index 2871f474bd..874dd6bdd5 100644 --- a/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java +++ b/common/transport/http/src/main/java/org/thingsboard/server/transport/http/DeviceApiController.java @@ -393,7 +393,8 @@ public class DeviceApiController implements TbTransportService { } @Override - public void onAttributeUpdate(AttributeUpdateNotificationMsg msg) { + public void onAttributeUpdate(UUID sessionId, AttributeUpdateNotificationMsg msg) { + log.trace("[{}] Received attributes update notification to device", sessionId); responseWriter.setResult(new ResponseEntity<>(JsonConverter.toJson(msg).toString(), HttpStatus.OK)); } @@ -404,7 +405,8 @@ public class DeviceApiController implements TbTransportService { } @Override - public void onToDeviceRpcRequest(ToDeviceRpcRequestMsg msg) { + public void onToDeviceRpcRequest(UUID sessionId, ToDeviceRpcRequestMsg msg) { + log.trace("[{}] Received RPC command to device", sessionId); responseWriter.setResult(new ResponseEntity<>(JsonConverter.toJson(msg, true).toString(), HttpStatus.OK)); transportService.process(sessionInfo, msg, false, TransportServiceCallback.EMPTY); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java index 8ffc3ac1f2..30ac8e01c3 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapConfig.java @@ -21,10 +21,11 @@ import org.eclipse.leshan.core.request.BindingMode; import org.eclipse.leshan.core.util.Hex; import org.eclipse.leshan.server.bootstrap.BootstrapConfig; +import java.io.Serializable; import java.nio.charset.StandardCharsets; @Data -public class LwM2MBootstrapConfig { +public class LwM2MBootstrapConfig implements Serializable { /* interface BootstrapSecurityConfig servers: BootstrapServersSecurityConfig, diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java index 2a2af5f099..bf30723ce0 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java @@ -146,10 +146,10 @@ public class LwM2mCredentialsSecurityInfoValidator { PSKClientCredentials pskConfig = (PSKClientCredentials) clientCredentialsConfig; if (StringUtils.isNotEmpty(pskConfig.getIdentity())) { try { - if (pskConfig.getKey() != null && pskConfig.getKey().length > 0) { + if (pskConfig.getDecodedKey() != null && pskConfig.getDecodedKey().length > 0) { endpoint = StringUtils.isNotEmpty(pskConfig.getEndpoint()) ? pskConfig.getEndpoint() : endpoint; if (endpoint != null && !endpoint.isEmpty()) { - result.setSecurityInfo(SecurityInfo.newPreSharedKeyInfo(endpoint, pskConfig.getIdentity(), pskConfig.getKey())); + result.setSecurityInfo(SecurityInfo.newPreSharedKeyInfo(endpoint, pskConfig.getIdentity(), pskConfig.getDecodedKey())); result.setSecurityMode(PSK); } } @@ -164,8 +164,8 @@ public class LwM2mCredentialsSecurityInfoValidator { private void createClientSecurityInfoRPK(TbLwM2MSecurityInfo result, String endpoint, LwM2MClientCredentials clientCredentialsConfig) { RPKClientCredentials rpkConfig = (RPKClientCredentials) clientCredentialsConfig; try { - if (rpkConfig.getKey() != null) { - PublicKey key = SecurityUtil.publicKey.decode(rpkConfig.getKey()); + if (rpkConfig.getDecodedKey() != null) { + PublicKey key = SecurityUtil.publicKey.decode(rpkConfig.getDecodedKey()); result.setSecurityInfo(SecurityInfo.newRawPublicKeyInfo(endpoint, key)); result.setSecurityMode(RPK); } else { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MDtlsCertificateVerifier.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MDtlsCertificateVerifier.java index 792ba131e8..04b69c815f 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MDtlsCertificateVerifier.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MDtlsCertificateVerifier.java @@ -42,8 +42,8 @@ import org.thingsboard.server.common.transport.util.SslUtil; import org.thingsboard.server.queue.util.TbLwM2mTransportComponent; import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; import org.thingsboard.server.transport.lwm2m.secure.credentials.LwM2MCredentials; -import org.thingsboard.server.transport.lwm2m.server.store.TbEditableSecurityStore; import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore; +import org.thingsboard.server.transport.lwm2m.server.store.TbMainSecurityStore; import javax.annotation.PostConstruct; import javax.security.auth.x500.X500Principal; @@ -67,7 +67,7 @@ public class TbLwM2MDtlsCertificateVerifier implements NewAdvancedCertificateVer private final TbLwM2MDtlsSessionStore sessionStorage; private final LwM2MTransportServerConfig config; private final LwM2mCredentialsSecurityInfoValidator securityInfoValidator; - private final TbEditableSecurityStore securityStore; + private final TbMainSecurityStore securityStore; @SuppressWarnings("deprecation") private StaticCertificateVerifier staticCertificateVerifier; @@ -134,7 +134,7 @@ public class TbLwM2MDtlsCertificateVerifier implements NewAdvancedCertificateVer if (msg.hasDeviceInfo() && deviceProfile != null) { sessionStorage.put(endpoint, new TbX509DtlsSessionInfo(cert.getSubjectX500Principal().getName(), msg)); try { - securityStore.put(securityInfo); + securityStore.putX509(securityInfo); } catch (NonUniqueSecurityInfoException e) { log.trace("Failed to add security info: {}", securityInfo, e); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MSecurityInfo.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MSecurityInfo.java index 9b9147c44f..bc45a77b58 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MSecurityInfo.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbLwM2MSecurityInfo.java @@ -23,8 +23,10 @@ import org.thingsboard.server.common.data.DeviceProfile; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2MBootstrapConfig; +import java.io.Serializable; + @Data -public class TbLwM2MSecurityInfo { +public class TbLwM2MSecurityInfo implements Serializable { private ValidateDeviceCredentialsResponse msg; private SecurityInfo securityInfo; private SecurityMode securityMode; diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbX509DtlsSessionInfo.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbX509DtlsSessionInfo.java index 1c038a9440..03490fe69d 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbX509DtlsSessionInfo.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/TbX509DtlsSessionInfo.java @@ -18,8 +18,10 @@ package org.thingsboard.server.transport.lwm2m.secure; import lombok.Data; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; +import java.io.Serializable; + @Data -public class TbX509DtlsSessionInfo { +public class TbX509DtlsSessionInfo implements Serializable { private final String x509CommonName; private final ValidateDeviceCredentialsResponse credentials; diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java index 36fff8e3d1..f5920cbddd 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java @@ -55,7 +55,8 @@ public class LwM2mSessionMsgListener implements GenericFutureListener lwM2mClientsByEndpoint = new ConcurrentHashMap<>(); private final Map lwM2mClientsByRegistrationId = new ConcurrentHashMap<>(); private final Map profiles = new ConcurrentHashMap<>(); @@ -75,6 +76,9 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { oldSession = lwM2MClient.getSession(); TbLwM2MSecurityInfo securityInfo = securityStore.getTbLwM2MSecurityInfoByEndpoint(lwM2MClient.getEndpoint()); if (securityInfo.getSecurityMode() != null) { + if (SecurityMode.X509.equals(securityInfo.getSecurityMode())) { + securityStore.registerX509(registration.getEndpoint(), registration.getId()); + } if (securityInfo.getDeviceProfile() != null) { profileUpdate(securityInfo.getDeviceProfile()); if (securityInfo.getSecurityInfo() != null) { @@ -124,7 +128,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext { if (currentRegistration.getId().equals(registration.getId())) { lwM2MClient.setState(LwM2MClientState.UNREGISTERED); lwM2mClientsByEndpoint.remove(lwM2MClient.getEndpoint()); - this.securityStore.remove(lwM2MClient.getEndpoint()); + this.securityStore.remove(lwM2MClient.getEndpoint(), registration.getId()); UUID profileId = lwM2MClient.getProfileId(); if (profileId != null) { Optional otherClients = lwM2mClientsByRegistrationId.values().stream().filter(e -> e.getProfileId().equals(profileId)).findFirst(); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MDtlsSessionRedisStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MDtlsSessionRedisStore.java index b1c4b85e2a..616688a620 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MDtlsSessionRedisStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2MDtlsSessionRedisStore.java @@ -15,28 +15,29 @@ */ package org.thingsboard.server.transport.lwm2m.server.store; -import com.fasterxml.jackson.databind.JsonNode; +import org.nustaq.serialization.FSTConfiguration; import org.springframework.data.redis.connection.RedisConnectionFactory; -import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.transport.lwm2m.secure.TbX509DtlsSessionInfo; public class TbLwM2MDtlsSessionRedisStore implements TbLwM2MDtlsSessionStore { private static final String SESSION_EP = "SESSION#EP#"; - RedisConnectionFactory connectionFactory; + private final RedisConnectionFactory connectionFactory; + private final FSTConfiguration serializer; public TbLwM2MDtlsSessionRedisStore(RedisConnectionFactory redisConnectionFactory) { this.connectionFactory = redisConnectionFactory; + this.serializer = FSTConfiguration.createDefaultConfiguration(); } @Override public void put(String endpoint, TbX509DtlsSessionInfo msg) { try (var c = connectionFactory.getConnection()) { - var msgJson = JacksonUtil.convertValue(msg, JsonNode.class); - if (msgJson != null) { - c.set(getKey(endpoint), msgJson.toString().getBytes()); + var serializedMsg = serializer.asByteArray(msg); + if (serializedMsg != null) { + c.set(getKey(endpoint), serializedMsg); } else { - throw new RuntimeException("Problem with serialization of message: " + msg.toString()); + throw new RuntimeException("Problem with serialization of message: " + msg); } } } @@ -46,7 +47,7 @@ public class TbLwM2MDtlsSessionRedisStore implements TbLwM2MDtlsSessionStore { try (var c = connectionFactory.getConnection()) { var data = c.get(getKey(endpoint)); if (data != null) { - return JacksonUtil.fromString(new String(data), TbX509DtlsSessionInfo.class); + return (TbX509DtlsSessionInfo) serializer.asObject(data); } else { return null; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java index 9e3fe5625d..b108286afd 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mRedisSecurityStore.java @@ -15,49 +15,55 @@ */ package org.thingsboard.server.transport.lwm2m.server.store; -import org.eclipse.leshan.server.redis.serialization.SecurityInfoSerDes; -import org.eclipse.leshan.server.security.EditableSecurityStore; import org.eclipse.leshan.server.security.NonUniqueSecurityInfoException; import org.eclipse.leshan.server.security.SecurityInfo; -import org.eclipse.leshan.server.security.SecurityStoreListener; -import org.springframework.data.redis.connection.RedisClusterConnection; +import org.nustaq.serialization.FSTConfiguration; import org.springframework.data.redis.connection.RedisConnectionFactory; -import org.springframework.data.redis.core.Cursor; -import org.springframework.data.redis.core.ScanOptions; +import org.springframework.integration.redis.util.RedisLockRegistry; import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MSecurityInfo; -import java.util.ArrayList; -import java.util.Collection; -import java.util.LinkedList; -import java.util.List; +import java.util.concurrent.locks.Lock; public class TbLwM2mRedisSecurityStore implements TbEditableSecurityStore { private static final String SEC_EP = "SEC#EP#"; - + private static final String LOCK_EP = "LOCK#EP#"; private static final String PSKID_SEC = "PSKID#SEC"; private final RedisConnectionFactory connectionFactory; - private SecurityStoreListener listener; + private final FSTConfiguration serializer; + private final RedisLockRegistry redisLock; public TbLwM2mRedisSecurityStore(RedisConnectionFactory connectionFactory) { this.connectionFactory = connectionFactory; + redisLock = new RedisLockRegistry(connectionFactory, "Security"); + serializer = FSTConfiguration.createDefaultConfiguration(); } @Override public SecurityInfo getByEndpoint(String endpoint) { + Lock lock = null; try (var connection = connectionFactory.getConnection()) { + lock = redisLock.obtain(toLockKey(endpoint)); + lock.lock(); byte[] data = connection.get((SEC_EP + endpoint).getBytes()); if (data == null) { return null; } else { - return deserialize(data); + return ((TbLwM2MSecurityInfo) serializer.asObject(data)).getSecurityInfo(); + } + } finally { + if (lock != null) { + lock.unlock(); } } } @Override public SecurityInfo getByIdentity(String identity) { + Lock lock = null; try (var connection = connectionFactory.getConnection()) { + lock = redisLock.obtain(toLockKey(identity)); + lock.lock(); byte[] ep = connection.hGet(PSKID_SEC.getBytes(), identity.getBytes()); if (ep == null) { return null; @@ -66,102 +72,86 @@ public class TbLwM2mRedisSecurityStore implements TbEditableSecurityStore { if (data == null) { return null; } else { - return deserialize(data); + return ((TbLwM2MSecurityInfo) serializer.asObject(data)).getSecurityInfo(); } } + } finally { + if (lock != null) { + lock.unlock(); + } } } @Override public void put(TbLwM2MSecurityInfo tbSecurityInfo) throws NonUniqueSecurityInfoException { - //TODO: implement + SecurityInfo info = tbSecurityInfo.getSecurityInfo(); + byte[] tbSecurityInfoSerialized = serializer.asByteArray(tbSecurityInfo); + Lock lock = null; + try (var connection = connectionFactory.getConnection()) { + lock = redisLock.obtain(tbSecurityInfo.getEndpoint()); + lock.lock(); + if (info != null && info.getIdentity() != null) { + byte[] oldEndpointBytes = connection.hGet(PSKID_SEC.getBytes(), info.getIdentity().getBytes()); + if (oldEndpointBytes != null) { + String oldEndpoint = new String(oldEndpointBytes); + if (!oldEndpoint.equals(info.getEndpoint())) { + throw new NonUniqueSecurityInfoException("PSK Identity " + info.getIdentity() + " is already used"); + } + connection.hSet(PSKID_SEC.getBytes(), info.getIdentity().getBytes(), info.getEndpoint().getBytes()); + } + } + + byte[] previousData = connection.getSet((SEC_EP + tbSecurityInfo.getEndpoint()).getBytes(), tbSecurityInfoSerialized); + if (previousData != null && info != null) { + String previousIdentity = ((TbLwM2MSecurityInfo) serializer.asObject(previousData)).getSecurityInfo().getIdentity(); + if (previousIdentity != null && !previousIdentity.equals(info.getIdentity())) { + connection.hDel(PSKID_SEC.getBytes(), previousIdentity.getBytes()); + } + } + } finally { + if (lock != null) { + lock.unlock(); + } + } } @Override public TbLwM2MSecurityInfo getTbLwM2MSecurityInfoByEndpoint(String endpoint) { - //TODO: implement - return null; + Lock lock = null; + try (var connection = connectionFactory.getConnection()) { + lock = redisLock.obtain(endpoint); + lock.lock(); + byte[] data = connection.get((SEC_EP + endpoint).getBytes()); + return (TbLwM2MSecurityInfo) serializer.asObject(data); + } finally { + if (lock != null) { + lock.unlock(); + } + } } @Override public void remove(String endpoint) { - //TODO: implement - } - - // @Override -// public Collection getAll() { -// try (var connection = connectionFactory.getConnection()) { -// Collection list = new LinkedList<>(); -// ScanOptions scanOptions = ScanOptions.scanOptions().count(100).match(SEC_EP + "*").build(); -// List> scans = new ArrayList<>(); -// if (connection instanceof RedisClusterConnection) { -// ((RedisClusterConnection) connection).clusterGetNodes().forEach(node -> { -// scans.add(((RedisClusterConnection) connection).scan(node, scanOptions)); -// }); -// } else { -// scans.add(connection.scan(scanOptions)); -// } -// -// scans.forEach(scan -> { -// scan.forEachRemaining(key -> { -// byte[] element = connection.get(key); -// list.add(deserialize(element)); -// }); -// }); -// return list; -// } -// } -// -// @Override -// public SecurityInfo add(SecurityInfo info) throws NonUniqueSecurityInfoException { -// byte[] data = serialize(info); -// try (var connection = connectionFactory.getConnection()) { -// if (info.getIdentity() != null) { -// // populate the secondary index (security info by PSK id) -// String oldEndpoint = new String(connection.hGet(PSKID_SEC.getBytes(), info.getIdentity().getBytes())); -// if (!oldEndpoint.equals(info.getEndpoint())) { -// throw new NonUniqueSecurityInfoException("PSK Identity " + info.getIdentity() + " is already used"); -// } -// connection.hSet(PSKID_SEC.getBytes(), info.getIdentity().getBytes(), info.getEndpoint().getBytes()); -// } -// -// byte[] previousData = connection.getSet((SEC_EP + info.getEndpoint()).getBytes(), data); -// SecurityInfo previous = previousData == null ? null : deserialize(previousData); -// String previousIdentity = previous == null ? null : previous.getIdentity(); -// if (previousIdentity != null && !previousIdentity.equals(info.getIdentity())) { -// connection.hDel(PSKID_SEC.getBytes(), previousIdentity.getBytes()); -// } -// -// return previous; -// } -// } -// -// @Override -// public SecurityInfo remove(String endpoint, boolean infosAreCompromised) { -// try (var connection = connectionFactory.getConnection()) { -// byte[] data = connection.get((SEC_EP + endpoint).getBytes()); -// -// if (data != null) { -// SecurityInfo info = deserialize(data); -// if (info.getIdentity() != null) { -// connection.hDel(PSKID_SEC.getBytes(), info.getIdentity().getBytes()); -// } -// connection.del((SEC_EP + endpoint).getBytes()); -// if (listener != null) { -// listener.securityInfoRemoved(infosAreCompromised, info); -// } -// return info; -// } -// } -// return null; -// } - - private byte[] serialize(SecurityInfo secInfo) { - return SecurityInfoSerDes.serialize(secInfo); + Lock lock = null; + try (var connection = connectionFactory.getConnection()) { + lock = redisLock.obtain(endpoint); + lock.lock(); + byte[] data = connection.get((SEC_EP + endpoint).getBytes()); + if (data != null) { + SecurityInfo info = ((TbLwM2MSecurityInfo) serializer.asObject(data)).getSecurityInfo(); + if (info != null && info.getIdentity() != null) { + connection.hDel(PSKID_SEC.getBytes(), info.getIdentity().getBytes()); + } + connection.del((SEC_EP + endpoint).getBytes()); + } + } finally { + if (lock != null) { + lock.unlock(); + } + } } - private SecurityInfo deserialize(byte[] data) { - return SecurityInfoSerDes.deserialize(data); + private String toLockKey(String endpoint) { + return LOCK_EP + endpoint; } - } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mSecurityStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mSecurityStore.java index d47be49978..bf1f275f32 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mSecurityStore.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mSecurityStore.java @@ -22,13 +22,22 @@ import org.jetbrains.annotations.Nullable; import org.thingsboard.server.transport.lwm2m.secure.LwM2mCredentialsSecurityInfoValidator; import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MSecurityInfo; +import java.util.HashSet; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.locks.Lock; +import java.util.concurrent.locks.ReentrantLock; + import static org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mTypeServer.CLIENT; @Slf4j -public class TbLwM2mSecurityStore implements TbEditableSecurityStore { +public class TbLwM2mSecurityStore implements TbMainSecurityStore { private final TbEditableSecurityStore securityStore; private final LwM2mCredentialsSecurityInfoValidator validator; + private final ConcurrentMap> endpointRegistrations = new ConcurrentHashMap<>(); public TbLwM2mSecurityStore(TbEditableSecurityStore securityStore, LwM2mCredentialsSecurityInfoValidator validator) { this.securityStore = securityStore; @@ -61,24 +70,42 @@ public class TbLwM2mSecurityStore implements TbEditableSecurityStore { @Nullable public SecurityInfo fetchAndPutSecurityInfo(String credentialsId) { TbLwM2MSecurityInfo securityInfo = validator.getEndpointSecurityInfoByCredentialsId(credentialsId, CLIENT); - try { - if (securityInfo != null) { + doPut(securityInfo); + return securityInfo != null ? securityInfo.getSecurityInfo() : null; + } + + private void doPut(TbLwM2MSecurityInfo securityInfo) { + if (securityInfo != null) { + try { securityStore.put(securityInfo); + } catch (NonUniqueSecurityInfoException e) { + log.trace("Failed to add security info: {}", securityInfo, e); } - } catch (NonUniqueSecurityInfoException e) { - log.trace("Failed to add security info: {}", securityInfo, e); } - return securityInfo != null ? securityInfo.getSecurityInfo() : null; } @Override - public void put(TbLwM2MSecurityInfo tbSecurityInfo) throws NonUniqueSecurityInfoException { - securityStore.put(tbSecurityInfo); + public void putX509(TbLwM2MSecurityInfo securityInfo) throws NonUniqueSecurityInfoException { + securityStore.put(securityInfo); } @Override - public void remove(String endpoint) { - //TODO: Make sure we delay removal of security store from endpoint due to reg/unreg race condition. -// securityStore.remove(endpoint); + public void registerX509(String endpoint, String registrationId) { + endpointRegistrations.computeIfAbsent(endpoint, ep -> new HashSet<>()).add(registrationId); + } + + @Override + public void remove(String endpoint, String registrationId) { + Set epRegistrationIds = endpointRegistrations.get(endpoint); + boolean shouldRemove; + if (epRegistrationIds == null) { + shouldRemove = true; + } else { + epRegistrationIds.remove(registrationId); + shouldRemove = epRegistrationIds.isEmpty(); + } + if (shouldRemove) { + securityStore.remove(endpoint); + } } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java index 154de636de..b9eb865df5 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbLwM2mStoreFactory.java @@ -51,7 +51,7 @@ public class TbLwM2mStoreFactory { } @Bean - private TbEditableSecurityStore securityStore() { + private TbMainSecurityStore securityStore() { return new TbLwM2mSecurityStore(redisConfiguration.isPresent() && useRedis ? new TbLwM2mRedisSecurityStore(redisConfiguration.get().redisConnectionFactory()) : new TbInMemorySecurityStore(), validator); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbMainSecurityStore.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbMainSecurityStore.java new file mode 100644 index 0000000000..f4394fb337 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/store/TbMainSecurityStore.java @@ -0,0 +1,29 @@ +/** + * Copyright © 2016-2021 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.transport.lwm2m.server.store; + +import org.eclipse.leshan.server.security.NonUniqueSecurityInfoException; +import org.thingsboard.server.transport.lwm2m.secure.TbLwM2MSecurityInfo; + +public interface TbMainSecurityStore extends TbSecurityStore { + + void putX509(TbLwM2MSecurityInfo tbSecurityInfo) throws NonUniqueSecurityInfoException; + + void registerX509(String endpoint, String registrationId); + + void remove(String endpoint, String registrationId); + +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index 58ccd5f8a0..3d073298de 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -796,7 +796,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } @Override - public void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg notification) { + public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg notification) { + log.trace("[{}] Received attributes update notification to device", sessionId); try { deviceSessionCtx.getPayloadAdaptor().convertToPublish(deviceSessionCtx, notification).ifPresent(deviceSessionCtx.getChannel()::writeAndFlush); } catch (Exception e) { @@ -811,7 +812,7 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement } @Override - public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { + public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) { log.trace("[{}] Received RPC command to device", sessionId); try { deviceSessionCtx.getPayloadAdaptor().convertToPublish(deviceSessionCtx, rpcRequest) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java index 55d36e0060..fb41093c35 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/GatewayDeviceSessionCtx.java @@ -84,7 +84,8 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple } @Override - public void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg notification) { + public void onAttributeUpdate(UUID sessionId, TransportProtos.AttributeUpdateNotificationMsg notification) { + log.trace("[{}] Received attributes update notification to device", sessionId); try { parent.getPayloadAdaptor().convertToGatewayPublish(this, getDeviceInfo().getDeviceName(), notification).ifPresent(parent::writeAndFlush); } catch (Exception e) { @@ -93,7 +94,8 @@ public class GatewayDeviceSessionCtx extends MqttDeviceAwareSessionContext imple } @Override - public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg request) { + public void onToDeviceRpcRequest(UUID sessionId, TransportProtos.ToDeviceRpcRequestMsg request) { + log.trace("[{}] Received RPC command to device", sessionId); try { parent.getPayloadAdaptor().convertToGatewayPublish(this, getDeviceInfo().getDeviceName(), request).ifPresent( payload -> { diff --git a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java index 65759f2eab..1fba12782a 100644 --- a/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java +++ b/common/transport/snmp/src/main/java/org/thingsboard/server/transport/snmp/session/DeviceSessionContext.java @@ -128,7 +128,8 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S } @Override - public void onAttributeUpdate(AttributeUpdateNotificationMsg attributeUpdateNotification) { + public void onAttributeUpdate(UUID sessionId, AttributeUpdateNotificationMsg attributeUpdateNotification) { + log.trace("[{}] Received attributes update notification to device", sessionId); snmpTransportContext.getSnmpTransportService().onAttributeUpdate(this, attributeUpdateNotification); } @@ -138,7 +139,8 @@ public class DeviceSessionContext extends DeviceAwareSessionContext implements S } @Override - public void onToDeviceRpcRequest(ToDeviceRpcRequestMsg toDeviceRequest) { + public void onToDeviceRpcRequest(UUID sessionId, ToDeviceRpcRequestMsg toDeviceRequest) { + log.trace("[{}] Received RPC command to device", sessionId); snmpTransportContext.getSnmpTransportService().onToDeviceRpcRequest(this, toDeviceRequest); snmpTransportContext.getTransportService().process(getSessionInfo(), toDeviceRequest, false, TransportServiceCallback.EMPTY); } diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java index 2eaf8dbf7c..644da7f4ec 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/SessionMsgListener.java @@ -36,11 +36,11 @@ public interface SessionMsgListener { void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse); - void onAttributeUpdate(AttributeUpdateNotificationMsg attributeUpdateNotification); + void onAttributeUpdate(UUID sessionId, AttributeUpdateNotificationMsg attributeUpdateNotification); void onRemoteSessionCloseCommand(UUID sessionId, SessionCloseNotificationProto sessionCloseNotification); - void onToDeviceRpcRequest(ToDeviceRpcRequestMsg toDeviceRequest); + void onToDeviceRpcRequest(UUID sessionId, ToDeviceRpcRequestMsg toDeviceRequest); void onToServerRpcResponse(ToServerRpcResponseMsg toServerResponse); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java index 243787fed7..a3826f5844 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportService.java @@ -20,7 +20,6 @@ import org.thingsboard.server.common.data.DeviceTransportType; import org.thingsboard.server.common.transport.auth.GetOrCreateDeviceFromGatewayResponse; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.common.transport.service.SessionMetaData; -import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.gen.transport.TransportProtos.ClaimDeviceMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.GetDeviceCredentialsRequestMsg; @@ -47,6 +46,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg; import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToRPCMsg; import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionInfoProto; +import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.TransportToDeviceActorMsg; @@ -109,7 +109,7 @@ public interface TransportService { void process(SessionInfoProto sessionInfo, ToServerRpcRequestMsg msg, TransportServiceCallback callback); - void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, boolean isFailedRpc, TransportServiceCallback callback); + void process(SessionInfoProto sessionInfo, ToDeviceRpcRequestMsg msg, boolean isFailedRpc, TransportServiceCallback callback); void process(SessionInfoProto sessionInfo, SubscriptionInfoProto msg, TransportServiceCallback callback); diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java index a436028714..f6ef357a93 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/TransportDeviceInfo.java @@ -22,8 +22,10 @@ import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.TenantId; +import java.io.Serializable; + @Data -public class TransportDeviceInfo { +public class TransportDeviceInfo implements Serializable { private TenantId tenantId; private CustomerId customerId; diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/ValidateDeviceCredentialsResponse.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/ValidateDeviceCredentialsResponse.java index e1324791b3..d54dce18fb 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/ValidateDeviceCredentialsResponse.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/auth/ValidateDeviceCredentialsResponse.java @@ -19,9 +19,11 @@ import lombok.Builder; import lombok.Data; import org.thingsboard.server.common.data.DeviceProfile; +import java.io.Serializable; + @Data @Builder -public class ValidateDeviceCredentialsResponse implements DeviceProfileAware { +public class ValidateDeviceCredentialsResponse implements DeviceProfileAware, Serializable { private final TransportDeviceInfo deviceInfo; private final DeviceProfile deviceProfile; diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java index 0ec2888c3e..99c2fd8f17 100644 --- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java +++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java @@ -776,7 +776,7 @@ public class DefaultTransportService implements TransportService { listener.onGetAttributesResponse(toSessionMsg.getGetAttributesResponse()); } if (toSessionMsg.hasAttributeUpdateNotification()) { - listener.onAttributeUpdate(toSessionMsg.getAttributeUpdateNotification()); + listener.onAttributeUpdate(sessionId, toSessionMsg.getAttributeUpdateNotification()); } if (toSessionMsg.hasSessionCloseNotification()) { listener.onRemoteSessionCloseCommand(sessionId, toSessionMsg.getSessionCloseNotification()); @@ -785,7 +785,7 @@ public class DefaultTransportService implements TransportService { listener.onToTransportUpdateCredentials(toSessionMsg.getToTransportUpdateCredentialsNotification()); } if (toSessionMsg.hasToDeviceRequest()) { - listener.onToDeviceRpcRequest(toSessionMsg.getToDeviceRequest()); + listener.onToDeviceRpcRequest(sessionId, toSessionMsg.getToDeviceRequest()); } if (toSessionMsg.hasToServerResponse()) { String requestId = sessionId + "-" + toSessionMsg.getToServerResponse().getRequestId(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java index d254c5ee9a..60ef98a801 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeDao.java @@ -176,4 +176,10 @@ public interface EdgeDao extends Dao { * @return the list of rule chain objects */ ListenableFuture> findEdgesByTenantIdAndDashboardId(UUID tenantId, UUID dashboardId); + + /** + * Executes stored procedure to cleanup old edge events. + * @param ttl the ttl for edge events in seconds + */ + void cleanupEvents(long ttl); } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java index 2ff7d94f43..0089853969 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/edge/EdgeServiceImpl.java @@ -627,6 +627,11 @@ public class EdgeServiceImpl extends AbstractEntityService implements EdgeServic return result.toString(); } + @Override + public void cleanupEvents(long ttl) { + edgeDao.cleanupEvents(ttl); + } + private List findEdgeRuleChains(TenantId tenantId, EdgeId edgeId) { List result = new ArrayList<>(); PageLink pageLink = new PageLink(DEFAULT_LIMIT); 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 91bfa954fd..2785df90df 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 @@ -131,6 +131,11 @@ public class BaseEventService implements EventService { } while (eventPageData.hasNext()); } + @Override + public void cleanupEvents(long ttl, long debugTtl) { + eventDao.cleanupEvents(ttl, debugTtl); + } + private DataValidator eventValidator = new DataValidator() { @Override 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 ba4e86c95e..ceacbffd50 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 @@ -102,4 +102,10 @@ public interface EventDao extends Dao { */ 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. + * @param otherEventsTtl the ttl for events in seconds + * @param debugEventsTtl the ttl for debug events in seconds + */ + void cleanupEvents(long otherEventsTtl, long debugEventsTtl); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java b/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java index 4431356690..fcd679383f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java @@ -15,11 +15,33 @@ */ package org.thingsboard.server.dao.sql; +import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; +import javax.sql.DataSource; +import java.sql.SQLException; +import java.sql.SQLWarning; +import java.sql.Statement; + +@Slf4j public abstract class JpaAbstractDaoListeningExecutorService { @Autowired protected JpaExecutorService service; + @Autowired + protected DataSource dataSource; + + protected void printWarnings(Statement statement) throws SQLException { + SQLWarning warnings = statement.getWarnings(); + if (warnings != null) { + log.debug("{}", warnings.getMessage()); + SQLWarning nextWarning = warnings.getNextWarning(); + while (nextWarning != null) { + log.debug("{}", nextWarning.getMessage()); + nextWarning = nextWarning.getNextWarning(); + } + } + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java index f17196fe93..2249de109f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java @@ -40,6 +40,10 @@ import org.thingsboard.server.dao.model.sql.EdgeInfoEntity; import org.thingsboard.server.dao.relation.RelationDao; import org.thingsboard.server.dao.sql.JpaAbstractSearchTextDao; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -194,6 +198,24 @@ public class JpaEdgeDao extends JpaAbstractSearchTextDao imple return transformFromRelationToEdge(tenantId, relations); } + @Override + public void cleanupEvents(long ttl) { + log.info("Going to cleanup old edge events using ttl: {}s", ttl); + try (Connection connection = dataSource.getConnection(); + PreparedStatement stmt = connection.prepareStatement("call cleanup_edge_events_by_ttl(?,?)")) { + stmt.setLong(1, ttl); + stmt.setLong(2, 0); + stmt.execute(); + printWarnings(stmt); + try (ResultSet resultSet = stmt.getResultSet()) { + resultSet.next(); + log.info("Total edge events removed by TTL: [{}]", resultSet.getLong(1)); + } + } catch (SQLException e) { + log.error("SQLException occurred during edge events TTL task execution ", e); + } + } + private ListenableFuture> transformFromRelationToEdge(UUID tenantId, ListenableFuture> relations) { return Futures.transformAsync(relations, input -> { List> edgeFutures = new ArrayList<>(input.size()); 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 44848ec515..ee09bb7a90 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 @@ -27,7 +27,6 @@ import org.thingsboard.server.common.data.Event; import org.thingsboard.server.common.data.event.DebugEvent; import org.thingsboard.server.common.data.event.ErrorEventFilter; 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; @@ -40,6 +39,10 @@ import org.thingsboard.server.dao.event.EventDao; import org.thingsboard.server.dao.model.sql.EventEntity; import org.thingsboard.server.dao.sql.JpaAbstractDao; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; import java.util.List; import java.util.Objects; import java.util.Optional; @@ -256,6 +259,25 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen return DaoUtil.convertDataList(latest); } + @Override + public void cleanupEvents(long otherEventsTtl, long debugEventsTtl) { + log.info("Going to cleanup old events using debug events ttl: {}s and other events ttl: {}s", debugEventsTtl, otherEventsTtl); + try (Connection connection = dataSource.getConnection(); + PreparedStatement stmt = connection.prepareStatement("call cleanup_events_by_ttl(?,?,?)")) { + stmt.setLong(1, otherEventsTtl); + stmt.setLong(2, debugEventsTtl); + stmt.setLong(3, 0); + stmt.execute(); + printWarnings(stmt); + try (ResultSet resultSet = stmt.getResultSet()){ + resultSet.next(); + log.info("Total events removed by TTL: [{}]", resultSet.getLong(1)); + } + } catch (SQLException e) { + log.error("SQLException occurred during events TTL task execution ", e); + } + } + public Optional save(EventEntity entity, boolean ifNotExists) { log.debug("Save event [{}] ", entity); if (entity.getTenantId() == null) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java index f6a6b56be5..eb12984753 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java @@ -25,9 +25,14 @@ import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; import javax.annotation.Nullable; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; import java.util.List; import java.util.Objects; import java.util.concurrent.TimeUnit; @@ -62,6 +67,24 @@ public abstract class AbstractSqlTimeseriesDao extends BaseAbstractSqlTimeseries @Value("${sql.ttl.ts.ts_key_value_ttl:0}") private long systemTtl; + public void cleanup(long systemTtl) { + log.info("Going to cleanup old timeseries data using ttl: {}s", systemTtl); + try (Connection connection = dataSource.getConnection(); + PreparedStatement stmt = connection.prepareStatement("call cleanup_timeseries_by_ttl(?,?,?)")) { + stmt.setObject(1, ModelConstants.NULL_UUID); + stmt.setLong(2, systemTtl); + stmt.setLong(3, 0); + stmt.execute(); + printWarnings(stmt); + try (ResultSet resultSet = stmt.getResultSet()) { + resultSet.next(); + log.info("Total telemetry removed stats by TTL for entities: [{}]", resultSet.getLong(1)); + } + } catch (SQLException e) { + log.error("SQLException occurred during timeseries TTL task execution ", e); + } + } + protected ListenableFuture> processFindAllAsync(TenantId tenantId, EntityId entityId, List queries) { List>> futures = queries .stream() diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java index c01d91ec61..c8241714f4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/hsql/JpaHsqlTimeseriesDao.java @@ -54,4 +54,9 @@ public class JpaHsqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDa return Futures.transform(tsQueue.add(entity), v -> dataPointDays, MoreExecutors.directExecutor()); } + @Override + public void cleanup(long systemTtl) { + + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java index 64c074fd40..c23e615548 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/psql/JpaPsqlTimeseriesDao.java @@ -35,6 +35,10 @@ import org.thingsboard.server.dao.timeseries.SqlTsPartitionDate; import org.thingsboard.server.dao.util.PsqlDao; import org.thingsboard.server.dao.util.SqlTsDao; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.ResultSet; +import java.sql.SQLException; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneOffset; @@ -62,6 +66,7 @@ public class JpaPsqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDa @Value("${sql.postgres.ts_key_value_partitioning:MONTHS}") private String partitioning; + @Override protected void init() { super.init(); @@ -93,6 +98,30 @@ public class JpaPsqlTimeseriesDao extends AbstractChunkedAggregationTimeseriesDa return Futures.transform(tsQueue.add(entity), v -> dataPointDays, MoreExecutors.directExecutor()); } + @Override + public void cleanup(long systemTtl) { + cleanupPartitions(systemTtl); + super.cleanup(systemTtl); + } + + private void cleanupPartitions(long systemTtl) { + log.info("Going to cleanup old timeseries data partitions using partition type: {} and ttl: {}s", partitioning, systemTtl); + try (Connection connection = dataSource.getConnection(); + PreparedStatement stmt = connection.prepareStatement("call drop_partitions_by_max_ttl(?,?,?)")) { + stmt.setString(1, partitioning); + stmt.setLong(2, systemTtl); + stmt.setLong(3, 0); + stmt.execute(); + printWarnings(stmt); + try (ResultSet resultSet = stmt.getResultSet()) { + resultSet.next(); + log.info("Total partitions removed by TTL: [{}]", resultSet.getLong(1)); + } + } catch (SQLException e) { + log.error("SQLException occurred during TTL task execution ", e); + } + } + private void savePartitionIfNotExist(long ts) { if (!tsFormat.equals(SqlTsPartitionDate.INDEFINITE) && ts >= 0) { LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java index 7f798ddb4b..31f3407fbf 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java @@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; +import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; import org.thingsboard.server.dao.model.sqlts.timescale.ts.TimescaleTsKvEntity; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; @@ -45,6 +46,9 @@ import org.thingsboard.server.dao.util.TimescaleDBTsDao; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; +import java.sql.CallableStatement; +import java.sql.SQLException; +import java.sql.Types; import java.util.*; import java.util.concurrent.CompletableFuture; import java.util.function.Function; @@ -156,6 +160,11 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements } } + @Override + public void cleanup(long systemTtl) { + super.cleanup(systemTtl); + } + private ListenableFuture> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { String strKey = query.getKey(); Integer keyId = getOrSaveKeyId(strKey); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index 701c67a648..fb15af723a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java @@ -126,6 +126,11 @@ public class BaseTimeseriesService implements TimeseriesService { return timeseriesLatestDao.findAllKeysByEntityIds(tenantId, entityIds); } + @Override + public void cleanup(long systemTtl) { + timeseriesDao.cleanup(systemTtl); + } + @Override public ListenableFuture save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry) { validate(entityId); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 240d5a0b88..ce653e2e6e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -288,6 +288,11 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD } } + @Override + public void cleanup(long systemTtl) { + //Cleanup by TTL is native for Cassandra + } + private ListenableFuture> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { long minPartition = toPartitionTs(query.getStartTs()); long maxPartition = toPartitionTs(query.getEndTs()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java index 3b3eb4ee0a..e9af5f0b75 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java @@ -38,4 +38,6 @@ public interface TimeseriesDao { ListenableFuture remove(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query); ListenableFuture removePartition(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query); + + void cleanup(long systemTtl); } diff --git a/dao/src/main/resources/sql/schema-ts-psql.sql b/dao/src/main/resources/sql/schema-ts-psql.sql index 5683cc0a17..2744ff5a07 100644 --- a/dao/src/main/resources/sql/schema-ts-psql.sql +++ b/dao/src/main/resources/sql/schema-ts-psql.sql @@ -38,17 +38,18 @@ CREATE OR REPLACE PROCEDURE drop_partitions_by_max_ttl(IN partition_type varchar LANGUAGE plpgsql AS $$ DECLARE - max_tenant_ttl bigint; - max_customer_ttl bigint; - max_ttl bigint; - date timestamp; - partition_by_max_ttl_date varchar; - partition_month varchar; - partition_day varchar; - partition_year varchar; - partition varchar; - partition_to_delete varchar; - + max_tenant_ttl bigint; + max_customer_ttl bigint; + max_ttl bigint; + date timestamp; + partition_by_max_ttl_date varchar; + partition_by_max_ttl_month varchar; + partition_by_max_ttl_day varchar; + partition_by_max_ttl_year varchar; + partition varchar; + partition_year integer; + partition_month integer; + partition_day integer; BEGIN SELECT max(attribute_kv.long_v) @@ -65,53 +66,138 @@ BEGIN if max_ttl IS NOT NULL AND max_ttl > 0 THEN date := to_timestamp(EXTRACT(EPOCH FROM current_timestamp) - max_ttl); partition_by_max_ttl_date := get_partition_by_max_ttl_date(partition_type, date); + RAISE NOTICE 'Date by max ttl: %', date; RAISE NOTICE 'Partition by max ttl: %', partition_by_max_ttl_date; IF partition_by_max_ttl_date IS NOT NULL THEN CASE WHEN partition_type = 'DAYS' THEN - partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); - partition_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); - partition_day := SPLIT_PART(partition_by_max_ttl_date, '_', 5); + partition_by_max_ttl_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); + partition_by_max_ttl_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); + partition_by_max_ttl_day := SPLIT_PART(partition_by_max_ttl_date, '_', 5); WHEN partition_type = 'MONTHS' THEN - partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); - partition_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); + partition_by_max_ttl_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); + partition_by_max_ttl_month := SPLIT_PART(partition_by_max_ttl_date, '_', 4); ELSE - partition_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); + partition_by_max_ttl_year := SPLIT_PART(partition_by_max_ttl_date, '_', 3); END CASE; - FOR partition IN SELECT tablename - FROM pg_tables - WHERE schemaname = 'public' - AND tablename like 'ts_kv_' || '%' - AND tablename != 'ts_kv_latest' - AND tablename != 'ts_kv_dictionary' - AND tablename != 'ts_kv_indefinite' - LOOP - IF partition != partition_by_max_ttl_date THEN - IF partition_year IS NOT NULL THEN - IF SPLIT_PART(partition, '_', 3)::integer < partition_year::integer THEN - partition_to_delete := partition; - ELSE - IF partition_month IS NOT NULL THEN - IF SPLIT_PART(partition, '_', 4)::integer < partition_month::integer THEN - partition_to_delete := partition; + IF partition_by_max_ttl_year IS NULL THEN + RAISE NOTICE 'Failed to remove partitions by max ttl date due to partition_by_max_ttl_year is null!'; + ELSE + IF partition_type = 'YEARS' THEN + FOR partition IN SELECT tablename + FROM pg_tables + WHERE schemaname = 'public' + AND tablename like 'ts_kv_' || '%' + AND tablename != 'ts_kv_latest' + AND tablename != 'ts_kv_dictionary' + AND tablename != 'ts_kv_indefinite' + AND tablename != partition_by_max_ttl_date + LOOP + partition_year := SPLIT_PART(partition, '_', 3)::integer; + IF partition_year < partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + END IF; + END LOOP; + ELSE + IF partition_type = 'MONTHS' THEN + IF partition_by_max_ttl_month IS NULL THEN + RAISE NOTICE 'Failed to remove months partitions by max ttl date due to partition_by_max_ttl_month is null!'; + ELSE + FOR partition IN SELECT tablename + FROM pg_tables + WHERE schemaname = 'public' + AND tablename like 'ts_kv_' || '%' + AND tablename != 'ts_kv_latest' + AND tablename != 'ts_kv_dictionary' + AND tablename != 'ts_kv_indefinite' + AND tablename != partition_by_max_ttl_date + LOOP + partition_year := SPLIT_PART(partition, '_', 3)::integer; + IF partition_year > partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; ELSE - IF partition_day IS NOT NULL THEN - IF SPLIT_PART(partition, '_', 5)::integer < partition_day::integer THEN - partition_to_delete := partition; + IF partition_year < partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + ELSE + partition_month := SPLIT_PART(partition, '_', 4)::integer; + IF partition_year = partition_by_max_ttl_year::integer THEN + IF partition_month >= partition_by_max_ttl_month::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + END IF; END IF; END IF; END IF; + END LOOP; + END IF; + ELSE + IF partition_type = 'DAYS' THEN + IF partition_by_max_ttl_month IS NULL THEN + RAISE NOTICE 'Failed to remove days partitions by max ttl date due to partition_by_max_ttl_month is null!'; + ELSE + IF partition_by_max_ttl_day IS NULL THEN + RAISE NOTICE 'Failed to remove days partitions by max ttl date due to partition_by_max_ttl_day is null!'; + ELSE + FOR partition IN SELECT tablename + FROM pg_tables + WHERE schemaname = 'public' + AND tablename like 'ts_kv_' || '%' + AND tablename != 'ts_kv_latest' + AND tablename != 'ts_kv_dictionary' + AND tablename != 'ts_kv_indefinite' + AND tablename != partition_by_max_ttl_date + LOOP + partition_year := SPLIT_PART(partition, '_', 3)::integer; + IF partition_year > partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + IF partition_year < partition_by_max_ttl_year::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + ELSE + partition_month := SPLIT_PART(partition, '_', 4)::integer; + IF partition_month > partition_by_max_ttl_month::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + IF partition_month < partition_by_max_ttl_month::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + ELSE + partition_day := SPLIT_PART(partition, '_', 5)::integer; + IF partition_day >= partition_by_max_ttl_day::integer THEN + RAISE NOTICE 'Skip iteration! Partition: % is valid!', partition; + CONTINUE; + ELSE + IF partition_day < partition_by_max_ttl_day::integer THEN + RAISE NOTICE 'Partition to delete by max ttl: %', partition; + EXECUTE format('DROP TABLE IF EXISTS %I', partition); + deleted := deleted + 1; + END IF; + END IF; + END IF; + END IF; + END IF; + END IF; + END LOOP; END IF; END IF; END IF; - IF partition_to_delete IS NOT NULL THEN - RAISE NOTICE 'Partition to delete by max ttl: %', partition_to_delete; - EXECUTE format('DROP TABLE IF EXISTS %I', partition_to_delete); - partition_to_delete := NULL; - deleted := deleted + 1; - END IF; END IF; - END LOOP; + END IF; + END IF; END IF; END IF; END @@ -127,8 +213,6 @@ BEGIN partition := 'ts_kv_' || to_char(date, 'yyyy') || '_' || to_char(date, 'MM'); WHEN partition_type = 'YEARS' THEN partition := 'ts_kv_' || to_char(date, 'yyyy'); - WHEN partition_type = 'INDEFINITE' THEN - partition := NULL; ELSE partition := NULL; END CASE; 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 2d0914c7d4..0c4b673d7f 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 @@ -2378,7 +2378,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable { public void setUserCredentialsEnabled(UserId userId, boolean userCredentialsEnabled) { restTemplate.postForLocation( - baseURL + "/api/user/{userId}/userCredentialsEnabled?serCredentialsEnabled={serCredentialsEnabled}", + baseURL + "/api/user/{userId}/userCredentialsEnabled?userCredentialsEnabled={userCredentialsEnabled}", null, userId.getId(), userCredentialsEnabled); diff --git a/ui-ngx/src/assets/locale/locale.constant-en_US.json b/ui-ngx/src/assets/locale/locale.constant-en_US.json index 2956e6c0f6..45bf0d750a 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -2369,7 +2369,7 @@ "delete-resource-text": "Be careful, after the confirmation the resource will become unrecoverable.", "delete-resource-title": "Are you sure you want to delete the resource '{{resourceTitle}}'?", "delete-resources-action-title": "Delete { count, plural, 1 {1 resource} other {# resources} }", - "delete-resources-text": "Be careful, after the confirmation all selected resources will be removed.", + "delete-resources-text": "Please note that the selected resources, even if they are used in device profiles, will be deleted.", "delete-resources-title": "Are you sure you want to delete { count, plural, 1 {1 resource} other {# resources} }?", "download": "Download resource", "drop-file": "Drop a resource file or click to select a file to upload.",