From 104ad0deb9c79f392e1c4234aca7983202aa21f3 Mon Sep 17 00:00:00 2001 From: nick Date: Thu, 17 Oct 2024 13:57:50 +0300 Subject: [PATCH 1/7] activity: Failed to send message to core --- .../service/DefaultTransportService.java | 22 +++++++++++-------- 1 file changed, 13 insertions(+), 9 deletions(-) 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 8e48061db6..43db4b7b23 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 @@ -45,6 +45,7 @@ import org.thingsboard.server.common.data.ResourceType; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.Tenant; import org.thingsboard.server.common.data.device.data.PowerMode; +import org.thingsboard.server.common.data.exception.TenantNotFoundException; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; @@ -1100,16 +1101,19 @@ public class DefaultTransportService extends TransportActivityManager implements } private void sendToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, UUID routingKey, TransportServiceCallback callback) { - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); - if (log.isTraceEnabled()) { - log.trace("[{}][{}] Pushing to topic {} message {}", tenantId, entityId, tpi.getFullTopicName(), msg); + try { + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); + if (log.isTraceEnabled()) { + log.trace("[{}][{}] Pushing to topic {} message {}", tenantId, entityId, tpi.getFullTopicName(), msg); + } + TransportTbQueueCallback transportTbQueueCallback = callback != null ? + new TransportTbQueueCallback(callback) : null; + tbCoreProducerStats.incrementTotal(); + StatsCallback wrappedCallback = new StatsCallback(transportTbQueueCallback, tbCoreProducerStats); + tbCoreMsgProducer.send(tpi, new TbProtoQueueMsg<>(routingKey, msg), wrappedCallback); + } catch (TenantNotFoundException e) { + log.trace("Failed to send message to core. Tenant with ID [{}] not found in the database. Message delivery aborted.", tenantId, e); } - - TransportTbQueueCallback transportTbQueueCallback = callback != null ? - new TransportTbQueueCallback(callback) : null; - tbCoreProducerStats.incrementTotal(); - StatsCallback wrappedCallback = new StatsCallback(transportTbQueueCallback, tbCoreProducerStats); - tbCoreMsgProducer.send(tpi, new TbProtoQueueMsg<>(routingKey, msg), wrappedCallback); } private void sendToRuleEngine(TenantId tenantId, DeviceId deviceId, CustomerId customerId, TransportProtos.SessionInfoProto sessionInfo, JsonObject json, From 12c35d903ad42837092ae459374e01a3521ee568 Mon Sep 17 00:00:00 2001 From: nick Date: Thu, 17 Oct 2024 16:46:16 +0300 Subject: [PATCH 2/7] reportActivity: add error to callback --- .../service/DefaultTransportService.java | 23 +++++++++++-------- 1 file changed, 14 insertions(+), 9 deletions(-) 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 43db4b7b23..c337a285e7 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 @@ -1101,19 +1101,24 @@ public class DefaultTransportService extends TransportActivityManager implements } private void sendToCore(TenantId tenantId, EntityId entityId, ToCoreMsg msg, UUID routingKey, TransportServiceCallback callback) { + TopicPartitionInfo tpi; try { - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); - if (log.isTraceEnabled()) { - log.trace("[{}][{}] Pushing to topic {} message {}", tenantId, entityId, tpi.getFullTopicName(), msg); - } - TransportTbQueueCallback transportTbQueueCallback = callback != null ? - new TransportTbQueueCallback(callback) : null; - tbCoreProducerStats.incrementTotal(); - StatsCallback wrappedCallback = new StatsCallback(transportTbQueueCallback, tbCoreProducerStats); - tbCoreMsgProducer.send(tpi, new TbProtoQueueMsg<>(routingKey, msg), wrappedCallback); + tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); } catch (TenantNotFoundException e) { log.trace("Failed to send message to core. Tenant with ID [{}] not found in the database. Message delivery aborted.", tenantId, e); + tpi = TopicPartitionInfo.builder().topic(e.getMessage()).build(); + if (callback != null) { + callback.onError(e); + } + } + if (log.isTraceEnabled()) { + log.trace("[{}][{}] Pushing to topic {} message {}", tenantId, entityId, tpi.getFullTopicName(), msg); } + TransportTbQueueCallback transportTbQueueCallback = callback != null ? + new TransportTbQueueCallback(callback) : null; + tbCoreProducerStats.incrementTotal(); + StatsCallback wrappedCallback = new StatsCallback(transportTbQueueCallback, tbCoreProducerStats); + tbCoreMsgProducer.send(tpi, new TbProtoQueueMsg<>(routingKey, msg), wrappedCallback); } private void sendToRuleEngine(TenantId tenantId, DeviceId deviceId, CustomerId customerId, TransportProtos.SessionInfoProto sessionInfo, JsonObject json, From 2d4a4ad9916d18b84cb01b0c83d3c2d3cd49feb9 Mon Sep 17 00:00:00 2001 From: nick Date: Thu, 17 Oct 2024 17:18:52 +0300 Subject: [PATCH 3/7] reportActivity: add topic empty --- .../common/transport/service/DefaultTransportService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 c337a285e7..edb0e5c44a 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 @@ -1106,7 +1106,7 @@ public class DefaultTransportService extends TransportActivityManager implements tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); } catch (TenantNotFoundException e) { log.trace("Failed to send message to core. Tenant with ID [{}] not found in the database. Message delivery aborted.", tenantId, e); - tpi = TopicPartitionInfo.builder().topic(e.getMessage()).build(); + tpi = TopicPartitionInfo.builder().topic("").build(); if (callback != null) { callback.onError(e); } From 2c0b7eb2b3dca386e5ac4b91c082dc71b3de7f61 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 18 Oct 2024 14:11:59 +0300 Subject: [PATCH 4/7] Deprecate all queue types except Kafka and in-memory --- .../DefaultNotificationCenter.java | 4 +- .../service/update/DeprecationService.java | 65 +++++++++++++++++++ .../notification/NotificationApiTest.java | 17 +++-- .../info/GeneralNotificationInfo.java | 36 ++++++++++ .../azure/servicebus/TbServiceBusAdmin.java | 1 + .../server/queue/pubsub/TbPubSubAdmin.java | 1 + .../queue/rabbitmq/TbRabbitMqAdmin.java | 1 + .../server/queue/sqs/TbAwsSqsAdmin.java | 1 + .../notification/DefaultNotifications.java | 8 +++ .../rule/engine/api/NotificationCenter.java | 3 +- 10 files changed, 131 insertions(+), 6 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/update/DeprecationService.java create mode 100644 common/data/src/main/java/org/thingsboard/server/common/data/notification/info/GeneralNotificationInfo.java diff --git a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java index 84d3b6c378..71e2b576ed 100644 --- a/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java +++ b/application/src/main/java/org/thingsboard/server/service/notification/DefaultNotificationCenter.java @@ -40,6 +40,7 @@ import org.thingsboard.server.common.data.notification.NotificationRequestConfig import org.thingsboard.server.common.data.notification.NotificationRequestStats; import org.thingsboard.server.common.data.notification.NotificationRequestStatus; import org.thingsboard.server.common.data.notification.NotificationStatus; +import org.thingsboard.server.common.data.notification.info.GeneralNotificationInfo; import org.thingsboard.server.common.data.notification.info.RuleOriginatedNotificationInfo; import org.thingsboard.server.common.data.notification.settings.NotificationSettings; import org.thingsboard.server.common.data.notification.settings.UserNotificationSettings; @@ -187,7 +188,7 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple } @Override - public void sendGeneralWebNotification(TenantId tenantId, UsersFilter recipients, NotificationTemplate template) { + public void sendGeneralWebNotification(TenantId tenantId, UsersFilter recipients, NotificationTemplate template, GeneralNotificationInfo info) { NotificationTarget target = new NotificationTarget(); target.setTenantId(tenantId); PlatformUsersNotificationTargetConfig targetConfig = new PlatformUsersNotificationTargetConfig(); @@ -198,6 +199,7 @@ public class DefaultNotificationCenter extends AbstractSubscriptionService imple .tenantId(tenantId) .template(template) .targets(List.of(EntityId.NULL_UUID)) // this is temporary and will be removed when 'create from scratch' functionality is implemented for recipients + .info(info) .status(NotificationRequestStatus.PROCESSING) .build(); try { diff --git a/application/src/main/java/org/thingsboard/server/service/update/DeprecationService.java b/application/src/main/java/org/thingsboard/server/service/update/DeprecationService.java new file mode 100644 index 0000000000..ddb8080bcd --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/update/DeprecationService.java @@ -0,0 +1,65 @@ +/** + * Copyright © 2016-2024 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.update; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Service; +import org.thingsboard.rule.engine.api.NotificationCenter; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.notification.info.GeneralNotificationInfo; +import org.thingsboard.server.common.data.notification.targets.platform.SystemAdministratorsFilter; +import org.thingsboard.server.dao.notification.DefaultNotifications; +import org.thingsboard.server.queue.util.AfterStartUp; + +import java.util.Map; + +@Service +@Slf4j +@RequiredArgsConstructor +public class DeprecationService { + + private final NotificationCenter notificationCenter; + + @Value("${queue.type}") + private String queueType; + + @AfterStartUp(order = Integer.MAX_VALUE) + public void checkDeprecation() { + checkQueueTypeDeprecation(); + } + + private void checkQueueTypeDeprecation() { + String queueTypeName; + switch (queueType) { + case "aws-sqs" -> queueTypeName = "AWS SQS"; + case "pubsub" -> queueTypeName = "PubSub"; + case "service-bus" -> queueTypeName = "Azure Service Bus"; + case "rabbitmq" -> queueTypeName = "RabbitMQ"; + default -> { + return; + } + } + + log.warn("WARNING: {} queue type is deprecated and will be removed in ThingsBoard 4.0. Please migrate to Apache Kafka", queueTypeName); + notificationCenter.sendGeneralWebNotification(TenantId.SYS_TENANT_ID, new SystemAdministratorsFilter(), + DefaultNotifications.queueTypeDeprecation.toTemplate(), new GeneralNotificationInfo(Map.of( + "queueType", queueTypeName + ))); + } + +} diff --git a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java index 8bae1feca9..c7a1aeb9e8 100644 --- a/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java +++ b/application/src/test/java/org/thingsboard/server/service/notification/NotificationApiTest.java @@ -56,6 +56,7 @@ import org.thingsboard.server.common.data.notification.NotificationRequestStatus import org.thingsboard.server.common.data.notification.NotificationType; import org.thingsboard.server.common.data.notification.info.AlarmCommentNotificationInfo; import org.thingsboard.server.common.data.notification.info.EntityActionNotificationInfo; +import org.thingsboard.server.common.data.notification.info.GeneralNotificationInfo; import org.thingsboard.server.common.data.notification.rule.trigger.config.AlarmCommentNotificationRuleTriggerConfig; import org.thingsboard.server.common.data.notification.settings.MobileAppNotificationDeliveryMethodConfig; import org.thingsboard.server.common.data.notification.settings.NotificationSettings; @@ -84,6 +85,7 @@ import org.thingsboard.server.common.data.notification.template.WebDeliveryMetho import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.dao.notification.DefaultNotifications; +import org.thingsboard.server.dao.notification.DefaultNotifications.DefaultNotification; import org.thingsboard.server.dao.service.DaoSqlTest; import org.thingsboard.server.service.notification.channels.MicrosoftTeamsNotificationChannel; import org.thingsboard.server.service.notification.channels.TeamsAdaptiveCard; @@ -746,14 +748,21 @@ public class NotificationApiTest extends AbstractNotificationApiTest { getAnotherWsClient().registerWaitForUpdate(); - DefaultNotifications.DefaultNotification expectedNotification = DefaultNotifications.maintenanceWork; + NotificationTemplate template = DefaultNotification.builder() + .name("Test") + .subject("Testing ${subjectVariable}") + .text("Testing ${bodyVariable}") + .build().toTemplate(); notificationCenter.sendGeneralWebNotification(TenantId.SYS_TENANT_ID, new SystemAdministratorsFilter(), - expectedNotification.toTemplate()); + template, new GeneralNotificationInfo(Map.of( + "subjectVariable", "subject", + "bodyVariable", "body" + ))); getAnotherWsClient().waitForUpdate(true); Notification notification = getAnotherWsClient().getLastDataUpdate().getUpdate(); - assertThat(notification.getSubject()).isEqualTo(expectedNotification.getSubject()); - assertThat(notification.getText()).isEqualTo(expectedNotification.getText()); + assertThat(notification.getSubject()).isEqualTo("Testing subject"); + assertThat(notification.getText()).isEqualTo("Testing body"); } @Test diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/GeneralNotificationInfo.java b/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/GeneralNotificationInfo.java new file mode 100644 index 0000000000..0fd1a8f014 --- /dev/null +++ b/common/data/src/main/java/org/thingsboard/server/common/data/notification/info/GeneralNotificationInfo.java @@ -0,0 +1,36 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.common.data.notification.info; + +import lombok.AllArgsConstructor; +import lombok.Data; +import lombok.NoArgsConstructor; + +import java.util.Map; + +@Data +@NoArgsConstructor +@AllArgsConstructor +public class GeneralNotificationInfo implements RuleOriginatedNotificationInfo { + + private Map data; + + @Override + public Map getTemplateData() { + return data; + } + +} diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java index 07561e7452..3f920778e6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/azure/servicebus/TbServiceBusAdmin.java @@ -31,6 +31,7 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; @Slf4j +@Deprecated(forRemoval = true, since = "3.9") // for removal in 4.0 public class TbServiceBusAdmin implements TbQueueAdmin { private final String MAX_SIZE = "maxSizeInMb"; private final String MESSAGE_TIME_TO_LIVE = "messageTimeToLiveInSec"; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java index 8d4b5abde0..bbae199c4a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/pubsub/TbPubSubAdmin.java @@ -37,6 +37,7 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; @Slf4j +@Deprecated(forRemoval = true, since = "3.9") // for removal in 4.0 public class TbPubSubAdmin implements TbQueueAdmin { private static final String ACK_DEADLINE = "ackDeadlineInSec"; private static final String MESSAGE_RETENTION = "messageRetentionInSec"; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java index 2d0cd5b414..aad3c1c8cc 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/rabbitmq/TbRabbitMqAdmin.java @@ -27,6 +27,7 @@ import java.util.Map; import java.util.concurrent.TimeoutException; @Slf4j +@Deprecated(forRemoval = true, since = "3.9") // for removal in 4.0 public class TbRabbitMqAdmin implements TbQueueAdmin { private final Channel channel; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java index 92def925d7..6aa1ec4dad 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsAdmin.java @@ -36,6 +36,7 @@ import java.util.function.Function; import java.util.stream.Collectors; @Slf4j +@Deprecated(forRemoval = true, since = "3.9") // for removal in 4.0 public class TbAwsSqsAdmin implements TbQueueAdmin { private final Map attributes; diff --git a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java index 1a4c3ffab2..65a5763c5a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java +++ b/dao/src/main/java/org/thingsboard/server/dao/notification/DefaultNotifications.java @@ -372,6 +372,14 @@ public class DefaultNotifications { .build()) .build(); + public static final DefaultNotification queueTypeDeprecation = DefaultNotification.builder() + .name("Queue type deprecation") + .type(NotificationType.GENERAL) + .subject("WARNING: ${queueType} deprecation") + .text("${queueType} queue type is deprecated and will be removed in ThingsBoard 4.0. Please migrate to Apache Kafka") + .icon("warning").color(RED_COLOR) + .build(); + private final NotificationTemplateService templateService; private final NotificationRuleService ruleService; diff --git a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java index d108a6981d..ce2ff08734 100644 --- a/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java +++ b/rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/NotificationCenter.java @@ -23,6 +23,7 @@ import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.notification.NotificationDeliveryMethod; import org.thingsboard.server.common.data.notification.NotificationRequest; import org.thingsboard.server.common.data.notification.NotificationRequestStats; +import org.thingsboard.server.common.data.notification.info.GeneralNotificationInfo; import org.thingsboard.server.common.data.notification.targets.platform.UsersFilter; import org.thingsboard.server.common.data.notification.template.NotificationTemplate; @@ -32,7 +33,7 @@ public interface NotificationCenter { NotificationRequest processNotificationRequest(TenantId tenantId, NotificationRequest notificationRequest, FutureCallback callback); - void sendGeneralWebNotification(TenantId tenantId, UsersFilter recipients, NotificationTemplate template); + void sendGeneralWebNotification(TenantId tenantId, UsersFilter recipients, NotificationTemplate template, GeneralNotificationInfo info); void deleteNotificationRequest(TenantId tenantId, NotificationRequestId notificationRequestId); From d77e557945040180f9a65262c0311a7e6ea69df8 Mon Sep 17 00:00:00 2001 From: nick Date: Fri, 18 Oct 2024 14:32:26 +0300 Subject: [PATCH 5/7] activity: add return if callback error --- .../common/transport/service/DefaultTransportService.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 edb0e5c44a..15073ba841 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 @@ -1104,12 +1104,12 @@ public class DefaultTransportService extends TransportActivityManager implements TopicPartitionInfo tpi; try { tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); - } catch (TenantNotFoundException e) { - log.trace("Failed to send message to core. Tenant with ID [{}] not found in the database. Message delivery aborted.", tenantId, e); - tpi = TopicPartitionInfo.builder().topic("").build(); + } catch (Exception e) { + log.trace("Failed to send message to core. Tenant with ID [{}], entityType [{}], entityId [{}], routingKey [{}], \nmsg [{}].\n Message delivery aborted.", tenantId, entityId.getEntityType(), entityId.toString(), routingKey, msg, e); if (callback != null) { callback.onError(e); } + return; } if (log.isTraceEnabled()) { log.trace("[{}][{}] Pushing to topic {} message {}", tenantId, entityId, tpi.getFullTopicName(), msg); From 737962ae075743d6696a10627ea70ed75116801a Mon Sep 17 00:00:00 2001 From: nick Date: Thu, 24 Oct 2024 13:17:12 +0300 Subject: [PATCH 6/7] Activity: comments1 --- .../common/transport/service/DefaultTransportService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 15073ba841..44f61270b4 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 @@ -1105,7 +1105,7 @@ public class DefaultTransportService extends TransportActivityManager implements try { tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, entityId); } catch (Exception e) { - log.trace("Failed to send message to core. Tenant with ID [{}], entityType [{}], entityId [{}], routingKey [{}], \nmsg [{}].\n Message delivery aborted.", tenantId, entityId.getEntityType(), entityId.toString(), routingKey, msg, e); + log.warn("Failed to send message to core. Tenant with ID [{}], routingKey [{}], msg [{}]. Message delivery aborted.", tenantId, routingKey, msg, e); if (callback != null) { callback.onError(e); } From 6581e48055df839aa50e91fbba9af5f72bb6e854 Mon Sep 17 00:00:00 2001 From: nick Date: Tue, 29 Oct 2024 12:26:12 +0200 Subject: [PATCH 7/7] tbel: fix bug validate decimal with <+> --- .../src/main/java/org/thingsboard/script/api/tbel/TbUtils.java | 2 +- .../test/java/org/thingsboard/script/api/tbel/TbUtilsTest.java | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java index 7f7b1f699f..6b321a3570 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbUtils.java @@ -1255,7 +1255,7 @@ public class TbUtils { if (str == null || str.isEmpty()) { return -1; } - return str.matches("-?\\d+(\\.\\d+)?") ? DEC_RADIX : -1; + return str.matches("[+-]?\\d+(\\.\\d+)?") ? DEC_RADIX : -1; } public static int isHexadecimal(String str) { diff --git a/common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbUtilsTest.java b/common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbUtilsTest.java index b7276edd8d..1c5a1a2e1a 100644 --- a/common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbUtilsTest.java +++ b/common/script/script-api/src/test/java/org/thingsboard/script/api/tbel/TbUtilsTest.java @@ -213,6 +213,7 @@ public class TbUtilsTest { Assertions.assertEquals((Integer) 0, TbUtils.parseInt("0")); Assertions.assertEquals((Integer) 0, TbUtils.parseInt("-0")); + Assertions.assertEquals((Integer) 0, TbUtils.parseInt("+0")); Assertions.assertEquals(java.util.Optional.of(473).get(), TbUtils.parseInt("473")); Assertions.assertEquals(java.util.Optional.of(-255).get(), TbUtils.parseInt("-0xFF")); Assertions.assertThrows(NumberFormatException.class, () -> TbUtils.parseInt("-0xFF123"));