From 2b95d6cbaed06d6c1f752abb4b728ceffa61ac23 Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Fri, 31 Jul 2020 15:06:47 +0300 Subject: [PATCH 01/14] changed logic to report activity & improvements for clear alarm node --- .../state/DefaultDeviceStateService.java | 8 ++--- .../engine/action/TbAbstractAlarmNode.java | 2 +- .../rule/engine/action/TbClearAlarmNode.java | 36 ++++++++++++++++--- 3 files changed, 36 insertions(+), 10 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index 8acca0b4c7..f8f0c6040c 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -197,15 +197,15 @@ public class DefaultDeviceStateService implements DeviceStateService { if (lastReportedActivity > 0 && lastReportedActivity > lastSavedActivity) { DeviceStateData stateData = getOrFetchDeviceStateData(deviceId); if (stateData != null) { - DeviceState state = stateData.getState(); - stateData.getState().setLastActivityTime(lastReportedActivity); - stateData.getMetaData().putValue("scope", SERVER_SCOPE); - pushRuleEngineMessage(stateData, ACTIVITY_EVENT); save(deviceId, LAST_ACTIVITY_TIME, lastReportedActivity); deviceLastSavedActivity.put(deviceId, lastReportedActivity); + DeviceState state = stateData.getState(); if (!state.isActive()) { state.setActive(true); save(deviceId, ACTIVITY_STATE, state.isActive()); + state.setLastActivityTime(lastReportedActivity); + stateData.getMetaData().putValue("scope", SERVER_SCOPE); + pushRuleEngineMessage(stateData, ACTIVITY_EVENT); } } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java index b3344105bb..766522d207 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java @@ -41,7 +41,7 @@ public abstract class TbAbstractAlarmNode processAlarm(TbContext ctx, TbMsg msg) { String alarmType = TbNodeUtils.processPattern(this.config.getAlarmType(), msg.getMetaData()); - ListenableFuture latest = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), msg.getOriginator(), alarmType); - return Futures.transformAsync(latest, a -> { - if (a != null && !a.getStatus().isCleared()) { - return clearAlarm(ctx, msg, a); + if (msg.getOriginator().getEntityType().equals(EntityType.ALARM)) { + return clearAlarmFromOriginator(ctx, msg); + } else { + ListenableFuture latest = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), msg.getOriginator(), alarmType); + return Futures.transformAsync(latest, a -> { + if (a != null && !a.getStatus().isCleared()) { + return clearAlarm(ctx, msg, a); + } + return Futures.immediateFuture(new AlarmResult(false, false, false, null)); + }, ctx.getDbCallbackExecutor()); + } + } + + private ListenableFuture clearAlarmFromOriginator(TbContext ctx, TbMsg msg) { + ListenableFuture alarmByIdAsync = ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), new AlarmId(msg.getOriginator().getId())); + return Futures.transformAsync(alarmByIdAsync, alarm -> { + if (alarm != null && !alarm.getStatus().isCleared()) { + long clearTs = System.currentTimeMillis(); + ListenableFuture clearAlarmFuture = ctx.getAlarmService().clearAlarm(ctx.getTenantId(), alarm.getId(), alarm.getDetails(), clearTs); + return Futures.transformAsync(clearAlarmFuture, cleared -> { + if (cleared) { + alarm.setClearTs(clearTs); + AlarmStatus oldStatus = alarm.getStatus(); + AlarmStatus newStatus = oldStatus.isAck() ? AlarmStatus.CLEARED_ACK : AlarmStatus.CLEARED_UNACK; + alarm.setStatus(newStatus); + return Futures.immediateFuture(new AlarmResult(false, false, true, alarm)); + } + return Futures.immediateFuture(new AlarmResult(false, false, false, alarm)); + }, ctx.getDbCallbackExecutor()); } return Futures.immediateFuture(new AlarmResult(false, false, false, null)); }, ctx.getDbCallbackExecutor()); From addf0ede2789be19a53a51f0816ce6c56fe42c92 Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Fri, 31 Jul 2020 16:29:49 +0300 Subject: [PATCH 02/14] added buildAlarmDetails for clearAlarmByAlarmOriginator method --- .../rule/engine/action/TbClearAlarmNode.java | 27 ++++++++++--------- 1 file changed, 15 insertions(+), 12 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java index a7d3c50e7a..d3c6429e34 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java @@ -58,7 +58,7 @@ public class TbClearAlarmNode extends TbAbstractAlarmNode processAlarm(TbContext ctx, TbMsg msg) { String alarmType = TbNodeUtils.processPattern(this.config.getAlarmType(), msg.getMetaData()); if (msg.getOriginator().getEntityType().equals(EntityType.ALARM)) { - return clearAlarmFromOriginator(ctx, msg); + return clearAlarmByAlarmOriginator(ctx, msg); } else { ListenableFuture latest = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), msg.getOriginator(), alarmType); return Futures.transformAsync(latest, a -> { @@ -70,21 +70,24 @@ public class TbClearAlarmNode extends TbAbstractAlarmNode clearAlarmFromOriginator(TbContext ctx, TbMsg msg) { + private ListenableFuture clearAlarmByAlarmOriginator(TbContext ctx, TbMsg msg) { ListenableFuture alarmByIdAsync = ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), new AlarmId(msg.getOriginator().getId())); return Futures.transformAsync(alarmByIdAsync, alarm -> { if (alarm != null && !alarm.getStatus().isCleared()) { - long clearTs = System.currentTimeMillis(); - ListenableFuture clearAlarmFuture = ctx.getAlarmService().clearAlarm(ctx.getTenantId(), alarm.getId(), alarm.getDetails(), clearTs); - return Futures.transformAsync(clearAlarmFuture, cleared -> { - if (cleared) { - alarm.setClearTs(clearTs); - AlarmStatus oldStatus = alarm.getStatus(); - AlarmStatus newStatus = oldStatus.isAck() ? AlarmStatus.CLEARED_ACK : AlarmStatus.CLEARED_UNACK; - alarm.setStatus(newStatus); + ctx.logJsEvalRequest(); + ListenableFuture asyncDetails = buildAlarmDetails(ctx, msg, alarm.getDetails()); + return Futures.transformAsync(asyncDetails, details -> { + ctx.logJsEvalRequest(); + long clearTs = System.currentTimeMillis(); + ListenableFuture clearAlarmFuture = ctx.getAlarmService().clearAlarm(ctx.getTenantId(), alarm.getId(), details, clearTs); + return Futures.transformAsync(clearAlarmFuture, cleared -> { + if (cleared) { + alarm.setClearTs(clearTs); + alarm.setDetails(details); + } + alarm.setStatus(alarm.getStatus().isAck() ? AlarmStatus.CLEARED_ACK : AlarmStatus.CLEARED_UNACK); return Futures.immediateFuture(new AlarmResult(false, false, true, alarm)); - } - return Futures.immediateFuture(new AlarmResult(false, false, false, alarm)); + }, ctx.getDbCallbackExecutor()); }, ctx.getDbCallbackExecutor()); } return Futures.immediateFuture(new AlarmResult(false, false, false, null)); From 158436b14147709c8c7544e1a4763687b10f3d8f Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Fri, 31 Jul 2020 16:49:24 +0300 Subject: [PATCH 03/14] fix typo --- .../server/service/state/DefaultDeviceStateService.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index f8f0c6040c..0025f80d71 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -200,10 +200,10 @@ public class DefaultDeviceStateService implements DeviceStateService { save(deviceId, LAST_ACTIVITY_TIME, lastReportedActivity); deviceLastSavedActivity.put(deviceId, lastReportedActivity); DeviceState state = stateData.getState(); + state.setLastActivityTime(lastReportedActivity); if (!state.isActive()) { state.setActive(true); save(deviceId, ACTIVITY_STATE, state.isActive()); - state.setLastActivityTime(lastReportedActivity); stateData.getMetaData().putValue("scope", SERVER_SCOPE); pushRuleEngineMessage(stateData, ACTIVITY_EVENT); } From 26e84f7b1834955648ab54acd4ea7bf15941282e Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Mon, 3 Aug 2020 13:54:40 +0300 Subject: [PATCH 04/14] refactoring & added alarmCanBeClearedWithAlarmOriginator test --- .../rule/engine/action/TbClearAlarmNode.java | 35 +++---------- .../rule/engine/action/TbAlarmNodeTest.java | 51 +++++++++++++++++++ 2 files changed, 57 insertions(+), 29 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java index d3c6429e34..30bf2d4d6e 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbClearAlarmNode.java @@ -57,38 +57,15 @@ public class TbClearAlarmNode extends TbAbstractAlarmNode processAlarm(TbContext ctx, TbMsg msg) { String alarmType = TbNodeUtils.processPattern(this.config.getAlarmType(), msg.getMetaData()); + ListenableFuture alarmFuture; if (msg.getOriginator().getEntityType().equals(EntityType.ALARM)) { - return clearAlarmByAlarmOriginator(ctx, msg); + alarmFuture = ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), new AlarmId(msg.getOriginator().getId())); } else { - ListenableFuture latest = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), msg.getOriginator(), alarmType); - return Futures.transformAsync(latest, a -> { - if (a != null && !a.getStatus().isCleared()) { - return clearAlarm(ctx, msg, a); - } - return Futures.immediateFuture(new AlarmResult(false, false, false, null)); - }, ctx.getDbCallbackExecutor()); + alarmFuture = ctx.getAlarmService().findLatestByOriginatorAndType(ctx.getTenantId(), msg.getOriginator(), alarmType); } - } - - private ListenableFuture clearAlarmByAlarmOriginator(TbContext ctx, TbMsg msg) { - ListenableFuture alarmByIdAsync = ctx.getAlarmService().findAlarmByIdAsync(ctx.getTenantId(), new AlarmId(msg.getOriginator().getId())); - return Futures.transformAsync(alarmByIdAsync, alarm -> { - if (alarm != null && !alarm.getStatus().isCleared()) { - ctx.logJsEvalRequest(); - ListenableFuture asyncDetails = buildAlarmDetails(ctx, msg, alarm.getDetails()); - return Futures.transformAsync(asyncDetails, details -> { - ctx.logJsEvalRequest(); - long clearTs = System.currentTimeMillis(); - ListenableFuture clearAlarmFuture = ctx.getAlarmService().clearAlarm(ctx.getTenantId(), alarm.getId(), details, clearTs); - return Futures.transformAsync(clearAlarmFuture, cleared -> { - if (cleared) { - alarm.setClearTs(clearTs); - alarm.setDetails(details); - } - alarm.setStatus(alarm.getStatus().isAck() ? AlarmStatus.CLEARED_ACK : AlarmStatus.CLEARED_UNACK); - return Futures.immediateFuture(new AlarmResult(false, false, true, alarm)); - }, ctx.getDbCallbackExecutor()); - }, ctx.getDbCallbackExecutor()); + return Futures.transformAsync(alarmFuture, a -> { + if (a != null && !a.getStatus().isCleared()) { + return clearAlarm(ctx, msg, a); } return Futures.immediateFuture(new AlarmResult(false, false, false, null)); }, ctx.getDbCallbackExecutor()); diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java index 61b83d647d..114967eb51 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/action/TbAlarmNodeTest.java @@ -35,6 +35,7 @@ import org.thingsboard.rule.engine.api.TbContext; import org.thingsboard.rule.engine.api.TbNodeConfiguration; import org.thingsboard.rule.engine.api.TbNodeException; import org.thingsboard.server.common.data.alarm.Alarm; +import org.thingsboard.server.common.data.id.AlarmId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.RuleChainId; @@ -95,6 +96,7 @@ public class TbAlarmNodeTest { private ListeningExecutor dbExecutor; private EntityId originator = new DeviceId(UUIDs.timeBased()); + private EntityId alarmOriginator = new AlarmId(UUIDs.timeBased()); private TenantId tenantId = new TenantId(UUIDs.timeBased()); private TbMsgMetaData metaData = new TbMsgMetaData(); private String rawJson = "{\"name\": \"Vit\", \"passed\": 5}"; @@ -325,6 +327,55 @@ public class TbAlarmNodeTest { assertEquals(expectedAlarm, actualAlarm); } + @Test + public void alarmCanBeClearedWithAlarmOriginator() throws ScriptException, IOException { + initWithClearAlarmScript(); + metaData.putValue("key", "value"); + TbMsg msg = TbMsg.newMsg( "USER", alarmOriginator, metaData, TbMsgDataType.JSON, rawJson, ruleChainId, ruleNodeId); + + long oldEndDate = System.currentTimeMillis(); + AlarmId id = new AlarmId(alarmOriginator.getId()); + Alarm activeAlarm = Alarm.builder().type("SomeType").tenantId(tenantId).originator(originator).status(ACTIVE_UNACK).severity(WARNING).endTs(oldEndDate).build(); + activeAlarm.setId(id); + + when(detailsJs.executeJsonAsync(msg)).thenReturn(Futures.immediateFuture(null)); + when(alarmService.findAlarmByIdAsync(tenantId, id)).thenReturn(Futures.immediateFuture(activeAlarm)); + when(alarmService.clearAlarm(eq(activeAlarm.getTenantId()), eq(activeAlarm.getId()), org.mockito.Mockito.any(JsonNode.class), anyLong())).thenReturn(Futures.immediateFuture(true)); +// doAnswer((Answer) invocationOnMock -> (Alarm) (invocationOnMock.getArguments())[0]).when(alarmService).createOrUpdateAlarm(activeAlarm); + + node.onMsg(ctx, msg); + + verify(ctx).tellNext(any(), eq("Cleared")); + + ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(TbMsg.class); + ArgumentCaptor typeCaptor = ArgumentCaptor.forClass(String.class); + ArgumentCaptor originatorCaptor = ArgumentCaptor.forClass(EntityId.class); + ArgumentCaptor metadataCaptor = ArgumentCaptor.forClass(TbMsgMetaData.class); + ArgumentCaptor dataCaptor = ArgumentCaptor.forClass(String.class); + verify(ctx).transformMsg(msgCaptor.capture(), typeCaptor.capture(), originatorCaptor.capture(), metadataCaptor.capture(), dataCaptor.capture()); + + assertEquals("ALARM", typeCaptor.getValue()); + assertEquals(alarmOriginator, originatorCaptor.getValue()); + assertEquals("value", metadataCaptor.getValue().getValue("key")); + assertEquals(Boolean.TRUE.toString(), metadataCaptor.getValue().getValue(IS_CLEARED_ALARM)); + assertNotSame(metaData, metadataCaptor.getValue()); + + Alarm actualAlarm = new ObjectMapper().readValue(dataCaptor.getValue().getBytes(), Alarm.class); + Alarm expectedAlarm = Alarm.builder() + .tenantId(tenantId) + .originator(originator) + .status(CLEARED_UNACK) + .severity(WARNING) + .propagate(false) + .type("SomeType") + .details(null) + .endTs(oldEndDate) + .build(); + expectedAlarm.setId(id); + + assertEquals(expectedAlarm, actualAlarm); + } + private void initWithCreateAlarmScript() { try { TbCreateAlarmNodeConfiguration config = new TbCreateAlarmNodeConfiguration(); From 45e3c2297117ae29f800a4376f49c4a721f1a907 Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Mon, 3 Aug 2020 13:57:45 +0300 Subject: [PATCH 05/14] fix typo --- .../org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java index 766522d207..b3344105bb 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/action/TbAbstractAlarmNode.java @@ -41,7 +41,7 @@ public abstract class TbAbstractAlarmNode Date: Thu, 18 Jun 2020 15:27:11 +0200 Subject: [PATCH 06/14] enable default credential provider chain for aws sqs --- application/src/main/resources/thingsboard.yml | 4 ++++ .../server/queue/sqs/TbAwsSqsAdmin.java | 15 ++++++++++----- .../server/queue/sqs/TbAwsSqsSettings.java | 3 +++ 3 files changed, 17 insertions(+), 5 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index cde049b99f..f4d76be201 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -612,6 +612,10 @@ queue: notifications: "${TB_QUEUE_KAFKA_NOTIFICATIONS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" js-executor: "${TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600}" aws_sqs: + # @see https://docs.aws.amazon.com/sdk-for-java/v1/developer-guide/java-dg-roles.html + # setting this to true, will ignore the access keys below and instead use the + # default credential provider chain, which includes instance profile credentials etc. + use_default_credential_provider_chain: "${TB_QUEUE_AWS_SQS_USE_DEFAULT_CREDENTIAL_PROVIDER_CHAIN:false}" access_key_id: "${TB_QUEUE_AWS_SQS_ACCESS_KEY_ID:YOUR_KEY}" secret_access_key: "${TB_QUEUE_AWS_SQS_SECRET_ACCESS_KEY:YOUR_SECRET}" region: "${TB_QUEUE_AWS_SQS_REGION:YOUR_REGION}" 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 8e293e6dff..7c13e32740 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 @@ -15,9 +15,7 @@ */ package org.thingsboard.server.queue.sqs; -import com.amazonaws.auth.AWSCredentials; -import com.amazonaws.auth.AWSStaticCredentialsProvider; -import com.amazonaws.auth.BasicAWSCredentials; +import com.amazonaws.auth.*; import com.amazonaws.services.sqs.AmazonSQS; import com.amazonaws.services.sqs.AmazonSQSClientBuilder; import com.amazonaws.services.sqs.model.CreateQueueRequest; @@ -37,9 +35,16 @@ public class TbAwsSqsAdmin implements TbQueueAdmin { public TbAwsSqsAdmin(TbAwsSqsSettings sqsSettings, Map attributes) { this.attributes = attributes; - AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); + AWSCredentialsProvider credentialsProvider; + if (sqsSettings.getUseDefaultCredentialProviderChain()) { + credentialsProvider = new DefaultAWSCredentialsProviderChain(); + } else { + AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); + credentialsProvider = new AWSStaticCredentialsProvider(awsCredentials); + } + sqsClient = AmazonSQSClientBuilder.standard() - .withCredentials(new AWSStaticCredentialsProvider(awsCredentials)) + .withCredentials(credentialsProvider) .withRegion(sqsSettings.getRegion()) .build(); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java index 922a2b1062..7a5c3332ad 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsSettings.java @@ -27,6 +27,9 @@ import org.springframework.stereotype.Component; @Data public class TbAwsSqsSettings { + @Value("${queue.aws_sqs.use_default_credential_provider_chain}") + private Boolean useDefaultCredentialProviderChain; + @Value("${queue.aws_sqs.access_key_id}") private String accessKeyId; From cc79b026035e6b9711fd142405b06bf1fcae88ad Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 3 Aug 2020 14:57:07 +0300 Subject: [PATCH 07/14] added default credentials provider chain for aws sqs consumer and producer --- .../src/main/resources/thingsboard.yml | 3 --- .../server/queue/sqs/TbAwsSqsAdmin.java | 6 +++++- .../queue/sqs/TbAwsSqsConsumerTemplate.java | 16 ++++++++++++---- .../queue/sqs/TbAwsSqsProducerTemplate.java | 19 +++++++++++-------- 4 files changed, 28 insertions(+), 16 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index f4d76be201..73db8c6735 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -612,9 +612,6 @@ queue: notifications: "${TB_QUEUE_KAFKA_NOTIFICATIONS_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" js-executor: "${TB_QUEUE_KAFKA_JE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:104857600}" aws_sqs: - # @see https://docs.aws.amazon.com/sdk-for-java/v1/developer-guide/java-dg-roles.html - # setting this to true, will ignore the access keys below and instead use the - # default credential provider chain, which includes instance profile credentials etc. use_default_credential_provider_chain: "${TB_QUEUE_AWS_SQS_USE_DEFAULT_CREDENTIAL_PROVIDER_CHAIN:false}" access_key_id: "${TB_QUEUE_AWS_SQS_ACCESS_KEY_ID:YOUR_KEY}" secret_access_key: "${TB_QUEUE_AWS_SQS_SECRET_ACCESS_KEY:YOUR_SECRET}" 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 7c13e32740..f99755a0af 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 @@ -15,7 +15,11 @@ */ package org.thingsboard.server.queue.sqs; -import com.amazonaws.auth.*; +import com.amazonaws.auth.AWSCredentials; +import com.amazonaws.auth.AWSCredentialsProvider; +import com.amazonaws.auth.AWSStaticCredentialsProvider; +import com.amazonaws.auth.BasicAWSCredentials; +import com.amazonaws.auth.DefaultAWSCredentialsProviderChain; import com.amazonaws.services.sqs.AmazonSQS; import com.amazonaws.services.sqs.AmazonSQSClientBuilder; import com.amazonaws.services.sqs.model.CreateQueueRequest; diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java index f4f279a02b..3e30123980 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsConsumerTemplate.java @@ -16,8 +16,10 @@ package org.thingsboard.server.queue.sqs; import com.amazonaws.auth.AWSCredentials; +import com.amazonaws.auth.AWSCredentialsProvider; import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.auth.BasicAWSCredentials; +import com.amazonaws.auth.DefaultAWSCredentialsProviderChain; import com.amazonaws.services.sqs.AmazonSQS; import com.amazonaws.services.sqs.AmazonSQSClientBuilder; import com.amazonaws.services.sqs.model.DeleteMessageBatchRequestEntry; @@ -67,13 +69,19 @@ public class TbAwsSqsConsumerTemplate extends AbstractPara this.decoder = decoder; this.sqsSettings = sqsSettings; - AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); - AWSStaticCredentialsProvider credProvider = new AWSStaticCredentialsProvider(awsCredentials); + AWSCredentialsProvider credentialsProvider; + if (sqsSettings.getUseDefaultCredentialProviderChain()) { + credentialsProvider = new DefaultAWSCredentialsProviderChain(); + } else { + AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); + credentialsProvider = new AWSStaticCredentialsProvider(awsCredentials); + } - this.sqsClient = AmazonSQSClientBuilder.standard() - .withCredentials(credProvider) + sqsClient = AmazonSQSClientBuilder.standard() + .withCredentials(credentialsProvider) .withRegion(sqsSettings.getRegion()) .build(); + } @Override diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java index 6110d08c5e..79bc6c803a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java @@ -16,8 +16,10 @@ package org.thingsboard.server.queue.sqs; import com.amazonaws.auth.AWSCredentials; +import com.amazonaws.auth.AWSCredentialsProvider; import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.auth.BasicAWSCredentials; +import com.amazonaws.auth.DefaultAWSCredentialsProviderChain; import com.amazonaws.services.sqs.AmazonSQS; import com.amazonaws.services.sqs.AmazonSQSClientBuilder; import com.amazonaws.services.sqs.model.SendMessageRequest; @@ -26,7 +28,6 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; -import com.google.common.util.concurrent.MoreExecutors; import com.google.gson.Gson; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -39,7 +40,6 @@ import org.thingsboard.server.queue.common.DefaultTbQueueMsg; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; @Slf4j public class TbAwsSqsProducerTemplate implements TbQueueProducer { @@ -54,15 +54,18 @@ public class TbAwsSqsProducerTemplate implements TbQueuePr this.admin = admin; this.defaultTopic = defaultTopic; - AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); - AWSStaticCredentialsProvider credProvider = new AWSStaticCredentialsProvider(awsCredentials); + AWSCredentialsProvider credentialsProvider; + if (sqsSettings.getUseDefaultCredentialProviderChain()) { + credentialsProvider = new DefaultAWSCredentialsProviderChain(); + } else { + AWSCredentials awsCredentials = new BasicAWSCredentials(sqsSettings.getAccessKeyId(), sqsSettings.getSecretAccessKey()); + credentialsProvider = new AWSStaticCredentialsProvider(awsCredentials); + } - this.sqsClient = AmazonSQSClientBuilder.standard() - .withCredentials(credProvider) + sqsClient = AmazonSQSClientBuilder.standard() + .withCredentials(credentialsProvider) .withRegion(sqsSettings.getRegion()) .build(); - - producerExecutor = MoreExecutors.listeningDecorator(Executors.newCachedThreadPool()); } @Override From 090c36999f97a6427f308770d1bd7cf75b751309 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Mon, 3 Aug 2020 15:47:18 +0300 Subject: [PATCH 08/14] refactored --- .../queue/sqs/TbAwsSqsProducerTemplate.java | 3 +++ msa/js-executor/package-lock.json | 23 ++++++++++++++----- 2 files changed, 20 insertions(+), 6 deletions(-) diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java index 79bc6c803a..3e508b09e6 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/sqs/TbAwsSqsProducerTemplate.java @@ -28,6 +28,7 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; +import com.google.common.util.concurrent.MoreExecutors; import com.google.gson.Gson; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; @@ -40,6 +41,7 @@ import org.thingsboard.server.queue.common.DefaultTbQueueMsg; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; @Slf4j public class TbAwsSqsProducerTemplate implements TbQueueProducer { @@ -66,6 +68,7 @@ public class TbAwsSqsProducerTemplate implements TbQueuePr .withCredentials(credentialsProvider) .withRegion(sqsSettings.getRegion()) .build(); + producerExecutor = MoreExecutors.listeningDecorator(Executors.newCachedThreadPool()); } @Override diff --git a/msa/js-executor/package-lock.json b/msa/js-executor/package-lock.json index 5636c4f586..47d7dba886 100644 --- a/msa/js-executor/package-lock.json +++ b/msa/js-executor/package-lock.json @@ -1872,12 +1872,14 @@ "balanced-match": { "version": "1.0.0", "bundled": true, - "dev": true + "dev": true, + "optional": true }, "brace-expansion": { "version": "1.1.11", "bundled": true, "dev": true, + "optional": true, "requires": { "balanced-match": "^1.0.0", "concat-map": "0.0.1" @@ -1892,17 +1894,20 @@ "code-point-at": { "version": "1.1.0", "bundled": true, - "dev": true + "dev": true, + "optional": true }, "concat-map": { "version": "0.0.1", "bundled": true, - "dev": true + "dev": true, + "optional": true }, "console-control-strings": { "version": "1.1.0", "bundled": true, - "dev": true + "dev": true, + "optional": true }, "core-util-is": { "version": "1.0.2", @@ -2019,7 +2024,8 @@ "inherits": { "version": "2.0.3", "bundled": true, - "dev": true + "dev": true, + "optional": true }, "ini": { "version": "1.3.5", @@ -2031,6 +2037,7 @@ "version": "1.0.0", "bundled": true, "dev": true, + "optional": true, "requires": { "number-is-nan": "^1.0.0" } @@ -2045,6 +2052,7 @@ "version": "3.0.4", "bundled": true, "dev": true, + "optional": true, "requires": { "brace-expansion": "^1.1.7" } @@ -2156,7 +2164,8 @@ "number-is-nan": { "version": "1.0.1", "bundled": true, - "dev": true + "dev": true, + "optional": true }, "object-assign": { "version": "4.1.1", @@ -2168,6 +2177,7 @@ "version": "1.4.0", "bundled": true, "dev": true, + "optional": true, "requires": { "wrappy": "1" } @@ -2289,6 +2299,7 @@ "version": "1.0.2", "bundled": true, "dev": true, + "optional": true, "requires": { "code-point-at": "^1.0.0", "is-fullwidth-code-point": "^1.0.0", From a8fe8b9e10a8a0d5bf18b47e1f529c60c671b472 Mon Sep 17 00:00:00 2001 From: Dmytro Shvaika Date: Wed, 5 Aug 2020 13:08:40 +0300 Subject: [PATCH 09/14] fix typo in ConditionalOnExpression annotation --- .../server/transport/mqtt/MqttSslHandlerProvider.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java index 4e3a39699a..0c872c351f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttSslHandlerProvider.java @@ -53,7 +53,7 @@ import java.util.concurrent.TimeUnit; */ @Slf4j @Component("MqttSslHandlerProvider") -@ConditionalOnExpression("'${transport.type:null}'=='null' || ('${transport.type}'=='local' && '${transport.http.enabled}'=='true')") +@ConditionalOnExpression("'${transport.type:null}'=='null' || ('${transport.type}'=='local' && '${transport.mqtt.enabled}'=='true')") @ConditionalOnProperty(prefix = "transport.mqtt.ssl", value = "enabled", havingValue = "true", matchIfMissing = false) public class MqttSslHandlerProvider { From 86102aea06735fcd7f7f5d306a55817eb2b23d34 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 4 Aug 2020 12:38:58 +0300 Subject: [PATCH 10/14] added other parameters for queue kafka --- application/src/main/resources/thingsboard.yml | 10 ++++++++++ .../server/queue/kafka/TbKafkaSettings.java | 5 ++++- 2 files changed, 14 insertions(+), 1 deletion(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 73db8c6735..2a58209560 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -605,6 +605,16 @@ queue: max_poll_records: "${TB_QUEUE_KAFKA_MAX_POLL_RECORDS:8192}" max_partition_fetch_bytes: "${TB_QUEUE_KAFKA_MAX_PARTITION_FETCH_BYTES:16777216}" fetch_max_bytes: "${TB_QUEUE_KAFKA_FETCH_MAX_BYTES:134217728}" + other: +# Properties for Confluent cloud +# - key: "ssl.endpoint.identification.algorithm" +# value: "https" +# - key: "sasl.mechanism" +# value: "PLAIN" +# - key: "sasl.jaas.config" +# value: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"CLUSTER_API_KEY\" password=\"CLUSTER_API_SECRET\";" +# - key: "security.protocol" +# value: "SASL_SSL" topic-properties: rule-engine: "${TB_QUEUE_KAFKA_RE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" core: "${TB_QUEUE_KAFKA_CORE_TOPIC_PROPERTIES:retention.ms:604800000;segment.bytes:26214400;retention.bytes:1048576000}" diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java index 659dd19bda..d9da969eb9 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java @@ -16,10 +16,12 @@ package org.thingsboard.server.queue.kafka; import lombok.Getter; +import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.producer.ProducerConfig; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; +import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.stereotype.Component; import java.util.List; @@ -30,6 +32,7 @@ import java.util.Properties; */ @Slf4j @ConditionalOnExpression("'${queue.type:null}'=='kafka'") +@ConfigurationProperties(prefix = "queue.kafka") @Component public class TbKafkaSettings { @@ -67,7 +70,7 @@ public class TbKafkaSettings { @Getter private int fetchMaxBytes; - @Value("${kafka.other:#{null}}") + @Setter private List other; public Properties toProps() { From b445473e3cf0ea2ee138dcdd85d7e74d3af6eed1 Mon Sep 17 00:00:00 2001 From: nordmif Date: Thu, 9 Jul 2020 14:22:38 +0300 Subject: [PATCH 11/14] added logging of MQTT payload errors --- .../server/transport/mqtt/adaptors/JsonMqttAdaptor.java | 1 + 1 file changed, 1 insertion(+) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index 45370a2c89..51a4c359f6 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -207,6 +207,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return new JsonParser().parse(payload); } catch (JsonSyntaxException ex) { + log.error(payload); throw new AdaptorException(ex); } } From 045c8e3acaaaa38f61bb75c1997ac1c112e0db55 Mon Sep 17 00:00:00 2001 From: nordmif Date: Thu, 9 Jul 2020 14:27:08 +0300 Subject: [PATCH 12/14] added logging of MQTT payload errors --- .../server/transport/mqtt/adaptors/JsonMqttAdaptor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index 51a4c359f6..755ebc9334 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -207,7 +207,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return new JsonParser().parse(payload); } catch (JsonSyntaxException ex) { - log.error(payload); + log.error("Payload is in incorrect format: " + payload); throw new AdaptorException(ex); } } From e5a6ddb03b0eba03d76dca08ede14cec071cee28 Mon Sep 17 00:00:00 2001 From: nordmif Date: Thu, 9 Jul 2020 14:56:09 +0300 Subject: [PATCH 13/14] added logging of MQTT payload errors --- .../server/transport/mqtt/adaptors/JsonMqttAdaptor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index 755ebc9334..6a28d9b47e 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -207,7 +207,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return new JsonParser().parse(payload); } catch (JsonSyntaxException ex) { - log.error("Payload is in incorrect format: " + payload); + log.error("Payload is in incorrect format: {}", payload); throw new AdaptorException(ex); } } From 52bed1748958ecf97bd33470875dc444cbcadb6d Mon Sep 17 00:00:00 2001 From: nordmif Date: Tue, 21 Jul 2020 15:03:45 +0300 Subject: [PATCH 14/14] changed logging level to WARN --- .../server/transport/mqtt/adaptors/JsonMqttAdaptor.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java index 6a28d9b47e..caa627451f 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/JsonMqttAdaptor.java @@ -207,7 +207,7 @@ public class JsonMqttAdaptor implements MqttTransportAdaptor { try { return new JsonParser().parse(payload); } catch (JsonSyntaxException ex) { - log.error("Payload is in incorrect format: {}", payload); + log.warn("Payload is in incorrect format: {}", payload); throw new AdaptorException(ex); } }