|
|
|
@ -21,10 +21,9 @@ import com.google.common.util.concurrent.ListenableFuture; |
|
|
|
import com.google.common.util.concurrent.SettableFuture; |
|
|
|
import com.google.gson.Gson; |
|
|
|
import com.google.gson.JsonObject; |
|
|
|
import groovy.lang.Tuple; |
|
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
import org.javatuples.Pair; |
|
|
|
import org.passay.Rule; |
|
|
|
import org.apache.commons.lang3.tuple.ImmutablePair; |
|
|
|
import org.apache.commons.lang3.tuple.Pair; |
|
|
|
import org.springframework.stereotype.Component; |
|
|
|
import org.thingsboard.rule.engine.api.msg.DeviceAttributesEventNotificationMsg; |
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
|
@ -37,7 +36,6 @@ import org.thingsboard.server.common.data.id.AssetId; |
|
|
|
import org.thingsboard.server.common.data.id.CustomerId; |
|
|
|
import org.thingsboard.server.common.data.id.DashboardId; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityViewId; |
|
|
|
import org.thingsboard.server.common.data.id.RuleChainId; |
|
|
|
@ -45,7 +43,6 @@ import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.id.UserId; |
|
|
|
import org.thingsboard.server.common.data.kv.AttributeKey; |
|
|
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.rule.RuleChain; |
|
|
|
import org.thingsboard.server.common.msg.TbMsg; |
|
|
|
import org.thingsboard.server.common.msg.TbMsgMetaData; |
|
|
|
import org.thingsboard.server.common.msg.queue.ServiceQueue; |
|
|
|
@ -145,9 +142,9 @@ public class TelemetryProcessor extends BaseProcessor { |
|
|
|
String defaultQueueName = deviceProfile.getDefaultQueueName(); |
|
|
|
queueName = defaultQueueName != null ? defaultQueueName : ServiceQueue.MAIN; |
|
|
|
} |
|
|
|
return new Pair<>(queueName, ruleChainId); |
|
|
|
return new ImmutablePair<>(queueName, ruleChainId); |
|
|
|
} else { |
|
|
|
return new Pair<>(ServiceQueue.MAIN, null); |
|
|
|
return new ImmutablePair<>(ServiceQueue.MAIN, null); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@ -157,8 +154,8 @@ public class TelemetryProcessor extends BaseProcessor { |
|
|
|
JsonObject json = JsonUtils.getJsonObject(tsKv.getKvList()); |
|
|
|
metaData.putValue("ts", tsKv.getTs() + ""); |
|
|
|
Pair<String, RuleChainId> defaultQueueAndRuleChain = getDefaultQueueNameAndRuleChainId(tenantId, entityId); |
|
|
|
String queueName = defaultQueueAndRuleChain.getValue0(); |
|
|
|
RuleChainId ruleChainId = defaultQueueAndRuleChain.getValue1(); |
|
|
|
String queueName = defaultQueueAndRuleChain.getKey(); |
|
|
|
RuleChainId ruleChainId = defaultQueueAndRuleChain.getValue(); |
|
|
|
TbMsg tbMsg = TbMsg.newMsg(queueName, SessionMsgType.POST_TELEMETRY_REQUEST.name(), entityId, metaData, gson.toJson(json), ruleChainId, null); |
|
|
|
tbClusterService.pushMsgToRuleEngine(tenantId, tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { |
|
|
|
@Override |
|
|
|
@ -180,8 +177,8 @@ public class TelemetryProcessor extends BaseProcessor { |
|
|
|
SettableFuture<Void> futureToSet = SettableFuture.create(); |
|
|
|
JsonObject json = JsonUtils.getJsonObject(msg.getKvList()); |
|
|
|
Pair<String, RuleChainId> defaultQueueAndRuleChain = getDefaultQueueNameAndRuleChainId(tenantId, entityId); |
|
|
|
String queueName = defaultQueueAndRuleChain.getValue0(); |
|
|
|
RuleChainId ruleChainId = defaultQueueAndRuleChain.getValue1(); |
|
|
|
String queueName = defaultQueueAndRuleChain.getKey(); |
|
|
|
RuleChainId ruleChainId = defaultQueueAndRuleChain.getValue(); |
|
|
|
TbMsg tbMsg = TbMsg.newMsg(queueName, SessionMsgType.POST_ATTRIBUTES_REQUEST.name(), entityId, metaData, gson.toJson(json), ruleChainId, null); |
|
|
|
tbClusterService.pushMsgToRuleEngine(tenantId, tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { |
|
|
|
@Override |
|
|
|
@ -207,8 +204,8 @@ public class TelemetryProcessor extends BaseProcessor { |
|
|
|
@Override |
|
|
|
public void onSuccess(@Nullable List<Void> voids) { |
|
|
|
Pair<String, RuleChainId> defaultQueueAndRuleChain = getDefaultQueueNameAndRuleChainId(tenantId, entityId); |
|
|
|
String queueName = defaultQueueAndRuleChain.getValue0(); |
|
|
|
RuleChainId ruleChainId = defaultQueueAndRuleChain.getValue1(); |
|
|
|
String queueName = defaultQueueAndRuleChain.getKey(); |
|
|
|
RuleChainId ruleChainId = defaultQueueAndRuleChain.getValue(); |
|
|
|
TbMsg tbMsg = TbMsg.newMsg(queueName, DataConstants.ATTRIBUTES_UPDATED, entityId, metaData, gson.toJson(json), ruleChainId, null); |
|
|
|
tbClusterService.pushMsgToRuleEngine(tenantId, tbMsg.getOriginator(), tbMsg, new TbQueueCallback() { |
|
|
|
@Override |
|
|
|
|