From 920f34b8611ce6127f5c2eb0caa11d3ba7a644a0 Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Fri, 7 Jul 2023 16:09:30 +0300 Subject: [PATCH 1/4] UI: Move form table style from data key to form.scss --- .../basic/common/data-key-row.component.html | 4 +- .../basic/common/data-key-row.component.scss | 22 ------ .../common/data-keys-panel.component.html | 24 +++---- .../common/data-keys-panel.component.scss | 52 ++------------ ui-ngx/src/form.scss | 67 +++++++++++++++++++ 5 files changed, 88 insertions(+), 81 deletions(-) diff --git a/ui-ngx/src/app/modules/home/components/widget/config/basic/common/data-key-row.component.html b/ui-ngx/src/app/modules/home/components/widget/config/basic/common/data-key-row.component.html index b639d08512..18e307787a 100644 --- a/ui-ngx/src/app/modules/home/components/widget/config/basic/common/data-key-row.component.html +++ b/ui-ngx/src/app/modules/home/components/widget/config/basic/common/data-key-row.component.html @@ -15,7 +15,7 @@ limitations under the License. --> -
+
{{ 'datakey.timeseries' | translate }} @@ -157,7 +157,7 @@
-
+
+ + + + + + + diff --git a/ui-ngx/src/app/shared/components/unit-input.component.scss b/ui-ngx/src/app/shared/components/unit-input.component.scss new file mode 100644 index 0000000000..25280e51ec --- /dev/null +++ b/ui-ngx/src/app/shared/components/unit-input.component.scss @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2023 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. + */ +.tb-autocomplete.tb-unit-input-autocomplete { + .mat-mdc-option { + border-bottom: none; + .mdc-list-item__primary-text { + flex: 1; + display: flex; + flex-direction: row; + gap: 8px; + .tb-unit-name, .tb-unit-symbol { + font-size: 14px; + font-weight: 400; + line-height: 20px; + letter-spacing: 0.2px; + } + .tb-unit-symbol { + color: rgba(0, 0, 0, 0.38); + min-width: 22px; + text-align: end; + b { + color: rgba(0, 0, 0, 0.87); + } + } + } + } +} diff --git a/ui-ngx/src/app/shared/components/unit-input.component.ts b/ui-ngx/src/app/shared/components/unit-input.component.ts new file mode 100644 index 0000000000..4b70c31cec --- /dev/null +++ b/ui-ngx/src/app/shared/components/unit-input.component.ts @@ -0,0 +1,148 @@ +/// +/// Copyright © 2016-2023 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. +/// + +import { Component, ElementRef, forwardRef, Input, OnInit, ViewChild, ViewEncapsulation } from '@angular/core'; +import { ControlValueAccessor, FormControl, NG_VALUE_ACCESSOR, UntypedFormBuilder } from '@angular/forms'; +import { Observable, of } from 'rxjs'; +import { searchUnits, Unit, unitBySymbol, units } from '@shared/models/unit.models'; +import { map, mergeMap, startWith, tap } from 'rxjs/operators'; +import { TranslateService } from '@ngx-translate/core'; + +@Component({ + selector: 'tb-unit-input', + templateUrl: './unit-input.component.html', + styleUrls: ['./unit-input.component.scss'], + providers: [ + { + provide: NG_VALUE_ACCESSOR, + useExisting: forwardRef(() => UnitInputComponent), + multi: true + } + ], + encapsulation: ViewEncapsulation.None +}) +export class UnitInputComponent implements ControlValueAccessor, OnInit { + + unitsFormControl: FormControl; + + modelValue: string | null; + + @Input() + disabled: boolean; + + @ViewChild('unitInput', {static: true}) unitInput: ElementRef; + + filteredUnits: Observable>; + + searchText = ''; + + private dirty = false; + + private translatedUnits: Array = units.map(u => ({symbol: u.symbol, + name: this.translate.instant(u.name), + tags: u.tags})); + + private propagateChange = (_val: any) => {}; + + constructor(private fb: UntypedFormBuilder, + private translate: TranslateService) { + } + + ngOnInit() { + this.unitsFormControl = this.fb.control('', []); + this.filteredUnits = this.unitsFormControl.valueChanges + .pipe( + tap(value => { + this.updateView(value); + }), + startWith(''), + map(value => (value as Unit)?.symbol ? (value as Unit).symbol : (value ? value as string : '')), + mergeMap(symbol => this.fetchUnits(symbol) ) + ); + } + + writeValue(symbol?: string): void { + this.searchText = ''; + this.modelValue = symbol; + let res: Unit | string = null; + if (symbol) { + const unit = unitBySymbol(symbol); + res = unit ? unit : symbol; + } + this.unitsFormControl.patchValue(res, {emitEvent: false}); + this.dirty = true; + } + + onFocus() { + if (this.dirty) { + this.unitsFormControl.updateValueAndValidity({onlySelf: true, emitEvent: true}); + this.dirty = false; + } + } + + updateView(value: Unit | string | null) { + const res: string = (value as Unit)?.symbol ? (value as Unit)?.symbol : (value as string); + if (this.modelValue !== res) { + this.modelValue = res; + this.propagateChange(this.modelValue); + } + } + + displayUnitFn(unit?: Unit | string): string | undefined { + if (unit) { + if ((unit as Unit).symbol) { + return (unit as Unit).symbol; + } else { + return unit as string; + } + } + return undefined; + } + + fetchUnits(searchText?: string): Observable> { + this.searchText = searchText; + const result = searchUnits(this.translatedUnits, searchText); + if (result.length) { + return of(result); + } else { + return of([]); + } + } + + registerOnChange(fn: any): void { + this.propagateChange = fn; + } + + registerOnTouched(fn: any): void { + } + + setDisabledState(isDisabled: boolean): void { + this.disabled = isDisabled; + if (this.disabled) { + this.unitsFormControl.disable({emitEvent: false}); + } else { + this.unitsFormControl.enable({emitEvent: false}); + } + } + + clear() { + this.unitsFormControl.patchValue(null, {emitEvent: true}); + setTimeout(() => { + this.unitInput.nativeElement.blur(); + this.unitInput.nativeElement.focus(); + }, 0); + } +} diff --git a/ui-ngx/src/app/shared/models/unit.models.ts b/ui-ngx/src/app/shared/models/unit.models.ts new file mode 100644 index 0000000000..7ac0d4db57 --- /dev/null +++ b/ui-ngx/src/app/shared/models/unit.models.ts @@ -0,0 +1,70 @@ +/// +/// Copyright © 2016-2023 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. +/// + +export interface Unit { + name: string; + symbol: string; + tags: string[]; +} + +export const units: Array = [ + { + name: 'unit.celsius', + symbol: '°C', + tags: ['temperature'] + }, + { + name: 'unit.kelvin', + symbol: 'K', + tags: ['temperature'] + }, + { + name: 'unit.fahrenheit', + symbol: '°F', + tags: ['temperature'] + }, + { + name: 'unit.percentage', + symbol: '%', + tags: ['percentage'] + }, + { + name: 'unit.second', + symbol: 's', + tags: ['time'] + }, + { + name: 'unit.minute', + symbol: 'min', + tags: ['time'] + }, + { + name: 'unit.hour', + symbol: 'h', + tags: ['time'] + } +]; + +export const unitBySymbol = (symbol: string): Unit => units.find(u => u.symbol === symbol); + +const searchUnitTags = (unit: Unit, searchText: string): boolean => + !!unit.tags.find(t => t.toUpperCase().includes(searchText.toUpperCase())); + +export const searchUnits = (_units: Array, searchText: string): Array => _units.filter( + u => u.symbol.toUpperCase().includes(searchText.toUpperCase()) || + u.name.toUpperCase().includes(searchText.toUpperCase()) || + searchUnitTags(u, searchText) +); diff --git a/ui-ngx/src/app/shared/pipe/highlight.pipe.ts b/ui-ngx/src/app/shared/pipe/highlight.pipe.ts index 0f8595fc3b..5d712c5b93 100644 --- a/ui-ngx/src/app/shared/pipe/highlight.pipe.ts +++ b/ui-ngx/src/app/shared/pipe/highlight.pipe.ts @@ -18,11 +18,10 @@ import { Pipe, PipeTransform } from '@angular/core'; @Pipe({ name: 'highlight' }) export class HighlightPipe implements PipeTransform { - transform(text: string, search): string { + transform(text: string, search: string, includes = false, flags = 'i'): string { const pattern = search .replace(/[\-\[\]\/\{\}\(\)\*\+\?\.\\\^\$\|]/g, '\\$&'); - const regex = new RegExp('^' + pattern, 'i'); - + const regex = new RegExp((!includes ? '^' : '') + pattern, flags); return search ? text.replace(regex, match => `${match}`) : text; } } diff --git a/ui-ngx/src/app/shared/shared.module.ts b/ui-ngx/src/app/shared/shared.module.ts index cf8b271171..f7c37e4761 100644 --- a/ui-ngx/src/app/shared/shared.module.ts +++ b/ui-ngx/src/app/shared/shared.module.ts @@ -194,6 +194,7 @@ import { ShortNumberPipe } from '@shared/pipe/short-number.pipe'; import { ToggleHeaderComponent, ToggleOption } from '@shared/components/toggle-header.component'; import { RuleChainSelectComponent } from '@shared/components/rule-chain/rule-chain-select.component'; import { ToggleSelectComponent } from '@shared/components/toggle-select.component'; +import { UnitInputComponent } from '@shared/components/unit-input.component'; export function MarkedOptionsFactory(markedOptionsService: MarkedOptionsService) { return markedOptionsService; @@ -367,6 +368,7 @@ export function MarkedOptionsFactory(markedOptionsService: MarkedOptionsService) ToggleHeaderComponent, ToggleOption, ToggleSelectComponent, + UnitInputComponent, RuleChainSelectComponent ], imports: [ @@ -597,6 +599,7 @@ export function MarkedOptionsFactory(markedOptionsService: MarkedOptionsService) ToggleHeaderComponent, ToggleOption, ToggleSelectComponent, + UnitInputComponent, RuleChainSelectComponent ] }) 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 c5ec1fca40..82f10b7f4c 100644 --- a/ui-ngx/src/assets/locale/locale.constant-en_US.json +++ b/ui-ngx/src/assets/locale/locale.constant-en_US.json @@ -3860,6 +3860,15 @@ "just-now": "Just now", "ago": "ago" }, + "unit": { + "celsius": "Celsius", + "kelvin": "Kelvin", + "fahrenheit": "Fahrenheit", + "percentage": "Percentage", + "second": "Second", + "minute": "Minute", + "hour": "Hour" + }, "user": { "user": "User", "users": "Users", diff --git a/ui-ngx/src/form.scss b/ui-ngx/src/form.scss index 0d66e09def..fbd816272b 100644 --- a/ui-ngx/src/form.scss +++ b/ui-ngx/src/form.scss @@ -177,9 +177,15 @@ opacity: 0; } } + &:not(.mat-mdc-form-field-has-icon-suffix) { + .mat-mdc-text-field-wrapper { + &.mdc-text-field--outlined, &:not(.mdc-text-field--outlined) { + padding-right: 12px; + } + } + } .mat-mdc-text-field-wrapper { &.mdc-text-field--outlined, &:not(.mdc-text-field--outlined) { - padding-right: 12px; padding-left: 12px; &:not(.mdc-text-field--focused):not(.mdc-text-field--disabled):not(:hover) { .mdc-notched-outline__leading, .mdc-notched-outline__trailing { @@ -233,7 +239,9 @@ } &.number { .mat-mdc-text-field-wrapper { - padding-right: 4px; + &.mdc-text-field--outlined, &:not(.mdc-text-field--outlined) { + padding-right: 4px; + } .mat-mdc-form-field-infix { input.mdc-text-field__input[type=number]::-webkit-inner-spin-button, input.mdc-text-field__input[type=number]::-webkit-outer-spin-button { From 788ac46f080fd655278067f5a5b186f7f8dbc915 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Tue, 11 Jul 2023 16:13:55 +0300 Subject: [PATCH 3/4] External rule nodes improvement: When force ack create new message to free memory after acknowledgment. Set request timeout for some external nodes. --- .../rule/engine/aws/sns/TbSnsNode.java | 9 +++++++-- .../rule/engine/aws/sqs/TbSqsNode.java | 10 +++++++--- .../engine/external/TbAbstractExternalNode.java | 5 ++++- .../rule/engine/gcp/pubsub/TbPubSubNode.java | 15 ++++++++++++++- .../rule/engine/kafka/TbKafkaNode.java | 10 +++++----- .../rule/engine/mail/TbSendEmailNode.java | 8 ++++---- .../thingsboard/rule/engine/mqtt/TbMqttNode.java | 8 ++++---- .../engine/notification/TbNotificationNode.java | 9 +++++---- .../rule/engine/notification/TbSlackNode.java | 6 +++--- .../rule/engine/rabbitmq/TbRabbitMqNode.java | 6 +++--- .../rule/engine/rest/TbRestApiCallNode.java | 4 ++-- .../rule/engine/sms/TbSendSmsNode.java | 10 +++++----- 12 files changed, 63 insertions(+), 37 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java index bad37922e7..bb93cb8d63 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.aws.sns; +import com.amazonaws.ClientConfiguration; import com.amazonaws.auth.AWSCredentials; import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.auth.BasicAWSCredentials; @@ -71,6 +72,9 @@ public class TbSnsNode extends TbAbstractExternalNode { this.snsClient = AmazonSNSClient.builder() .withCredentials(credProvider) .withRegion(this.config.getRegion()) + .withClientConfiguration(new ClientConfiguration() + .withConnectionTimeout(10000) + .withRequestTimeout(5000)) .build(); } catch (Exception e) { throw new TbNodeException(e); @@ -79,9 +83,10 @@ public class TbSnsNode extends TbAbstractExternalNode { @Override public void onMsg(TbContext ctx, TbMsg msg) throws ExecutionException, InterruptedException, TbNodeException { - withCallback(publishMessageAsync(ctx, msg), + var tbMsg = ackIfNeeded(ctx, msg); + withCallback(publishMessageAsync(ctx, tbMsg), m -> tellSuccess(ctx, m), - t -> tellFailure(ctx, processException(ctx, msg, t), t)); + t -> tellFailure(ctx, processException(ctx, tbMsg, t), t)); ackIfNeeded(ctx, msg); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java index f8ebc8e295..2b9f8916cf 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sqs/TbSqsNode.java @@ -15,6 +15,7 @@ */ package org.thingsboard.rule.engine.aws.sqs; +import com.amazonaws.ClientConfiguration; import com.amazonaws.auth.AWSCredentials; import com.amazonaws.auth.AWSStaticCredentialsProvider; import com.amazonaws.auth.BasicAWSCredentials; @@ -78,6 +79,9 @@ public class TbSqsNode extends TbAbstractExternalNode { this.sqsClient = AmazonSQSClientBuilder.standard() .withCredentials(credProvider) .withRegion(this.config.getRegion()) + .withClientConfiguration(new ClientConfiguration() + .withConnectionTimeout(10000) + .withRequestTimeout(5000)) .build(); } catch (Exception e) { throw new TbNodeException(e); @@ -86,10 +90,10 @@ public class TbSqsNode extends TbAbstractExternalNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { - withCallback(publishMessageAsync(ctx, msg), + var tbMsg = ackIfNeeded(ctx, msg); + withCallback(publishMessageAsync(ctx, tbMsg), m -> tellSuccess(ctx, m), - t -> tellFailure(ctx, processException(ctx, msg, t), t)); - ackIfNeeded(ctx, msg); + t -> tellFailure(ctx, processException(ctx, tbMsg, t), t)); } private ListenableFuture publishMessageAsync(TbContext ctx, TbMsg msg) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/external/TbAbstractExternalNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/external/TbAbstractExternalNode.java index f5f8c34ab3..834dbe2390 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/external/TbAbstractExternalNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/external/TbAbstractExternalNode.java @@ -52,9 +52,12 @@ public abstract class TbAbstractExternalNode implements TbNode { } } - protected void ackIfNeeded(TbContext ctx, TbMsg msg) { + protected TbMsg ackIfNeeded(TbContext ctx, TbMsg msg) { if (forceAck) { ctx.ack(msg); + return msg.copyWithNewCtx(); + } else { + return msg; } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java index 3b927f05a6..18d817fd94 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/gcp/pubsub/TbPubSubNode.java @@ -20,6 +20,7 @@ import com.google.api.core.ApiFutureCallback; import com.google.api.core.ApiFutures; import com.google.api.gax.core.CredentialsProvider; import com.google.api.gax.core.FixedCredentialsProvider; +import com.google.api.gax.retrying.RetrySettings; import com.google.auth.oauth2.ServiceAccountCredentials; import com.google.cloud.pubsub.v1.Publisher; import com.google.protobuf.ByteString; @@ -36,6 +37,7 @@ import org.thingsboard.rule.engine.external.TbAbstractExternalNode; import org.thingsboard.server.common.data.plugin.ComponentType; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgMetaData; +import org.threeten.bp.Duration; import java.io.ByteArrayInputStream; import java.io.IOException; @@ -75,8 +77,8 @@ public class TbPubSubNode extends TbAbstractExternalNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { + msg = ackIfNeeded(ctx, msg); publishMessage(ctx, msg); - ackIfNeeded(ctx, msg); } @Override @@ -134,8 +136,19 @@ public class TbPubSubNode extends TbAbstractExternalNode { new ByteArrayInputStream(config.getServiceAccountKey().getBytes())); CredentialsProvider credProvider = FixedCredentialsProvider.create(credentials); + var retrySettings = RetrySettings.newBuilder() + .setTotalTimeout(Duration.ofSeconds(10)) + .setInitialRetryDelay(Duration.ofMillis(50)) + .setRetryDelayMultiplier(1.1) + .setMaxRetryDelay(Duration.ofSeconds(2)) + .setInitialRpcTimeout(Duration.ofSeconds(2)) + .setRpcTimeoutMultiplier(1) + .setMaxRpcTimeout(Duration.ofSeconds(10)) + .build(); + return Publisher.newBuilder(topicName) .setCredentialsProvider(credProvider) + .setRetrySettings(retrySettings) .build(); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java index ea8e3edb9d..760eec2421 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/kafka/TbKafkaNode.java @@ -115,25 +115,25 @@ public class TbKafkaNode extends TbAbstractExternalNode { public void onMsg(TbContext ctx, TbMsg msg) { String topic = TbNodeUtils.processPattern(config.getTopicPattern(), msg); String keyPattern = config.getKeyPattern(); + var tbMsg = ackIfNeeded(ctx, msg); try { if (initError != null) { - ctx.tellFailure(msg, new RuntimeException("Failed to initialize Kafka rule node producer: " + initError.getMessage())); + ctx.tellFailure(tbMsg, new RuntimeException("Failed to initialize Kafka rule node producer: " + initError.getMessage())); } else { ctx.getExternalCallExecutor().executeAsync(() -> { publish( ctx, - msg, + tbMsg, topic, keyPattern == null || keyPattern.isEmpty() ? null - : TbNodeUtils.processPattern(config.getKeyPattern(), msg) + : TbNodeUtils.processPattern(config.getKeyPattern(), tbMsg) ); return null; }); } - ackIfNeeded(ctx, msg); } catch (Exception e) { - ctx.tellFailure(msg, e); + ctx.tellFailure(tbMsg, e); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java index 3afea4f817..3e705d9f08 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mail/TbSendEmailNode.java @@ -75,13 +75,13 @@ public class TbSendEmailNode extends TbAbstractExternalNode { try { validateType(msg.getType()); TbEmail email = getEmail(msg); + var tbMsg = ackIfNeeded(ctx, msg); withCallback(ctx.getMailExecutor().executeAsync(() -> { - sendEmail(ctx, msg, email); + sendEmail(ctx, tbMsg, email); return null; }), - ok -> tellSuccess(ctx, msg), - fail -> tellFailure(ctx, msg, fail)); - ackIfNeeded(ctx, msg); + ok -> tellSuccess(ctx, tbMsg), + fail -> tellFailure(ctx, tbMsg, fail)); } catch (Exception ex) { ctx.tellFailure(msg, ex); } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java index 8fac7b1683..121b9fb756 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/mqtt/TbMqttNode.java @@ -80,16 +80,16 @@ public class TbMqttNode extends TbAbstractExternalNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { String topic = TbNodeUtils.processPattern(this.mqttNodeConfiguration.getTopicPattern(), msg); - this.mqttClient.publish(topic, Unpooled.wrappedBuffer(msg.getData().getBytes(UTF8)), MqttQoS.AT_LEAST_ONCE, mqttNodeConfiguration.isRetainedMessage()) + var tbMsg = ackIfNeeded(ctx, msg); + this.mqttClient.publish(topic, Unpooled.wrappedBuffer(tbMsg.getData().getBytes(UTF8)), MqttQoS.AT_LEAST_ONCE, mqttNodeConfiguration.isRetainedMessage()) .addListener(future -> { if (future.isSuccess()) { - tellSuccess(ctx, msg); + tellSuccess(ctx, tbMsg); } else { - tellFailure(ctx, processException(ctx, msg, future.cause()), future.cause()); + tellFailure(ctx, processException(ctx, tbMsg, future.cause()), future.cause()); } } ); - ackIfNeeded(ctx, msg); } private TbMsg processException(TbContext ctx, TbMsg origMsg, Throwable e) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java index bf57628d64..f1fd9727fa 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbNotificationNode.java @@ -71,16 +71,17 @@ public class TbNotificationNode extends TbAbstractExternalNode { .originatorEntityId(ctx.getSelf().getRuleChainId()) .build(); + var tbMsg = ackIfNeeded(ctx, msg); + DonAsynchron.withCallback(ctx.getNotificationExecutor().executeAsync(() -> ctx.getNotificationCenter().processNotificationRequest(ctx.getTenantId(), notificationRequest, stats -> { - TbMsgMetaData metaData = msg.getMetaData().copy(); + TbMsgMetaData metaData = tbMsg.getMetaData().copy(); metaData.putValue("notificationRequestResult", JacksonUtil.toString(stats)); - tellSuccess(ctx, TbMsg.transformMsg(msg, metaData)); + tellSuccess(ctx, TbMsg.transformMsg(tbMsg, metaData)); })), r -> { }, - e -> tellFailure(ctx, msg, e)); - ackIfNeeded(ctx, msg); + e -> tellFailure(ctx, tbMsg, e)); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java index fd56043848..86e2c4bd1c 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/notification/TbSlackNode.java @@ -61,12 +61,12 @@ public class TbSlackNode extends TbAbstractExternalNode { } String message = TbNodeUtils.processPattern(config.getMessageTemplate(), msg); + var tbMsg = ackIfNeeded(ctx, msg); DonAsynchron.withCallback(ctx.getExternalCallExecutor().executeAsync(() -> { ctx.getSlackService().sendMessage(ctx.getTenantId(), token, config.getConversation().getId(), message); }), - r -> tellSuccess(ctx, msg), - e -> tellFailure(ctx, msg, e)); - ackIfNeeded(ctx, msg); + r -> tellSuccess(ctx, tbMsg), + e -> tellFailure(ctx, tbMsg, e)); } } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java index 4ae634902f..4b42ee1d18 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rabbitmq/TbRabbitMqNode.java @@ -84,10 +84,10 @@ public class TbRabbitMqNode extends TbAbstractExternalNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { - withCallback(publishMessageAsync(ctx, msg), + var tbMsg = ackIfNeeded(ctx, msg); + withCallback(publishMessageAsync(ctx, tbMsg), m -> tellSuccess(ctx, m), - t -> tellFailure(ctx, processException(ctx, msg, t), t)); - ackIfNeeded(ctx, msg); + t -> tellFailure(ctx, processException(ctx, tbMsg, t), t)); } private ListenableFuture publishMessageAsync(TbContext ctx, TbMsg msg) { diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java index 94b0e5d078..c6083f1098 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rest/TbRestApiCallNode.java @@ -60,10 +60,10 @@ public class TbRestApiCallNode extends TbAbstractExternalNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { - httpClient.processMessage(ctx, msg, + var tbMsg = ackIfNeeded(ctx, msg); + httpClient.processMessage(ctx, tbMsg, m -> tellSuccess(ctx, m), (m, t) -> tellFailure(ctx, m, t)); - ackIfNeeded(ctx, msg); } @Override diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/sms/TbSendSmsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/sms/TbSendSmsNode.java index f6cbccd4d4..55209e1f64 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/sms/TbSendSmsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/sms/TbSendSmsNode.java @@ -60,16 +60,16 @@ public class TbSendSmsNode extends TbAbstractExternalNode { @Override public void onMsg(TbContext ctx, TbMsg msg) { + var tbMsg = ackIfNeeded(ctx, msg); try { withCallback(ctx.getSmsExecutor().executeAsync(() -> { - sendSms(ctx, msg); + sendSms(ctx, tbMsg); return null; }), - ok -> tellSuccess(ctx, msg), - fail -> tellFailure(ctx, msg, fail)); - ackIfNeeded(ctx, msg); + ok -> tellSuccess(ctx, tbMsg), + fail -> tellFailure(ctx, tbMsg, fail)); } catch (Exception ex) { - ctx.tellFailure(msg, ex); + ctx.tellFailure(tbMsg, ex); } } From 3c2d5eeb30fb7e249d52ce04a8924c190af1fbc6 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Tue, 11 Jul 2023 16:15:51 +0300 Subject: [PATCH 4/4] Minor fix --- .../main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java | 1 - 1 file changed, 1 deletion(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java index bb93cb8d63..a8a1137a98 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/aws/sns/TbSnsNode.java @@ -87,7 +87,6 @@ public class TbSnsNode extends TbAbstractExternalNode { withCallback(publishMessageAsync(ctx, tbMsg), m -> tellSuccess(ctx, m), t -> tellFailure(ctx, processException(ctx, tbMsg, t), t)); - ackIfNeeded(ctx, msg); } private ListenableFuture publishMessageAsync(TbContext ctx, TbMsg msg) {