Browse Source

added sequence for the all RPC

pull/5091/head
YevhenBondarenko 5 years ago
parent
commit
8513c99903
  1. 73
      application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java
  2. 5
      application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java
  4. 17
      application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java
  5. 5
      common/cluster-api/src/main/proto/queue.proto
  6. 2
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  7. 1
      common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java
  8. 25
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java
  9. 29
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java
  10. 24
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java
  11. 1
      rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java
  12. 6
      rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java

73
application/src/main/java/org/thingsboard/server/actors/device/DeviceActorMessageProcessor.java

@ -80,9 +80,9 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionType;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg; import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToAttributeUpdatesMsg;
import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToRPCMsg; import org.thingsboard.server.gen.transport.TransportProtos.SubscribeToRPCMsg;
import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionInfoProto; import org.thingsboard.server.gen.transport.TransportProtos.SubscriptionInfoProto;
import org.thingsboard.server.gen.transport.TransportProtos.ToDevicePersistedRpcResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcResponseStatusMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto; import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto;
@ -298,7 +298,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
} }
systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(),
null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION)); null, requestMd.isSent() ? RpcError.TIMEOUT : RpcError.NO_ACTIVE_CONNECTION));
sendNextPendingRequest(context); if (!requestMd.isDelivered()) {
sendNextPendingRequest(context);
}
} }
} }
@ -315,22 +317,10 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
} }
Set<Integer> sentOneWayIds = new HashSet<>(); Set<Integer> sentOneWayIds = new HashSet<>();
if (sessionType == SessionType.ASYNC) { if (rpcSequenceEnabled) {
if (rpcSequenceEnabled) { toDeviceRpcPendingMap.entrySet().stream().filter(e -> !e.getValue().isDelivered()).findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds));
List<Map.Entry<Integer, ToDeviceRpcRequestMetadata>> entries = new ArrayList<>(); } else if (sessionType == SessionType.ASYNC) {
for (Map.Entry<Integer, ToDeviceRpcRequestMetadata> entry : toDeviceRpcPendingMap.entrySet()) { toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds));
if (entry.getValue().isDelivered()) {
continue;
}
entries.add(entry);
if (entry.getValue().getMsg().getMsg().isPersisted() || entry.getValue().getMsg().getMsg().isOneway()) {
break;
}
}
entries.forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds));
} else {
toDeviceRpcPendingMap.entrySet().forEach(processPendingRpc(context, sessionId, nodeId, sentOneWayIds));
}
} else { } else {
toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds)); toDeviceRpcPendingMap.entrySet().stream().findFirst().ifPresent(processPendingRpc(context, sessionId, nodeId, sentOneWayIds));
} }
@ -348,7 +338,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
return entry -> { return entry -> {
ToDeviceRpcRequest request = entry.getValue().getMsg().getMsg(); ToDeviceRpcRequest request = entry.getValue().getMsg().getMsg();
ToDeviceRpcRequestBody body = request.getBody(); ToDeviceRpcRequestBody body = request.getBody();
if (request.isOneway()) { if (request.isOneway() && !rpcSequenceEnabled) {
sentOneWayIds.add(entry.getKey()); sentOneWayIds.add(entry.getKey());
systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(request.getId(), null, null)); systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(request.getId(), null, null));
} }
@ -357,6 +347,7 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
.setMethodName(body.getMethod()) .setMethodName(body.getMethod())
.setParams(body.getParams()) .setParams(body.getParams())
.setExpirationTime(request.getExpirationTime()) .setExpirationTime(request.getExpirationTime())
.setTimeout(request.getTimeout())
.setRequestIdMSB(request.getId().getMostSignificantBits()) .setRequestIdMSB(request.getId().getMostSignificantBits())
.setRequestIdLSB(request.getId().getLeastSignificantBits()) .setRequestIdLSB(request.getId().getLeastSignificantBits())
.setOneway(request.isOneway()) .setOneway(request.isOneway())
@ -395,8 +386,8 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
if (msg.hasClaimDevice()) { if (msg.hasClaimDevice()) {
handleClaimDeviceMsg(context, sessionInfo, msg.getClaimDevice()); handleClaimDeviceMsg(context, sessionInfo, msg.getClaimDevice());
} }
if (msg.hasPersistedRpcResponseMsg()) { if (msg.hasRpcResponseStatusMsg()) {
processPersistedRpcResponses(context, sessionInfo, msg.getPersistedRpcResponseMsg()); processPersistedRpcResponses(context, sessionInfo, msg.getRpcResponseStatusMsg());
} }
if (msg.hasUplinkNotificationMsg()) { if (msg.hasUplinkNotificationMsg()) {
processUplinkNotificationMsg(context, sessionInfo, msg.getUplinkNotificationMsg()); processUplinkNotificationMsg(context, sessionInfo, msg.getUplinkNotificationMsg());
@ -556,27 +547,32 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
boolean success = requestMd != null; boolean success = requestMd != null;
if (success) { if (success) {
boolean hasError = StringUtils.isNotEmpty(responseMsg.getError()); boolean hasError = StringUtils.isNotEmpty(responseMsg.getError());
String payload = hasError ? responseMsg.getError() : responseMsg.getPayload(); try {
systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor( String payload = hasError ? responseMsg.getError() : responseMsg.getPayload();
new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(), systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(
payload, null)); new FromDeviceRpcResponse(requestMd.getMsg().getMsg().getId(),
if (requestMd.getMsg().getMsg().isPersisted()) { payload, null));
RpcStatus status = hasError ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL; if (requestMd.getMsg().getMsg().isPersisted()) {
JsonNode response; RpcStatus status = hasError ? RpcStatus.FAILED : RpcStatus.SUCCESSFUL;
try { JsonNode response;
response = JacksonUtil.toJsonNode(payload); try {
} catch (IllegalArgumentException e) { response = JacksonUtil.toJsonNode(payload);
response = JacksonUtil.newObjectNode().put("error", payload); } catch (IllegalArgumentException e) {
response = JacksonUtil.newObjectNode().put("error", payload);
}
systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response);
}
} finally {
if (!requestMd.isDelivered() && hasError) {
sendNextPendingRequest(context);
} }
systemContext.getTbRpcService().save(tenantId, new RpcId(requestMd.getMsg().getMsg().getId()), status, response);
} }
sendNextPendingRequest(context);
} else { } else {
log.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId()); log.debug("[{}] Rpc command response [{}] is stale!", deviceId, responseMsg.getRequestId());
} }
} }
private void processPersistedRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDevicePersistedRpcResponseMsg responseMsg) { private void processPersistedRpcResponses(TbActorCtx context, SessionInfoProto sessionInfo, ToDeviceRpcResponseStatusMsg responseMsg) {
UUID rpcId = new UUID(responseMsg.getRequestIdMSB(), responseMsg.getRequestIdLSB()); UUID rpcId = new UUID(responseMsg.getRequestIdMSB(), responseMsg.getRequestIdLSB());
RpcStatus status = RpcStatus.valueOf(responseMsg.getStatus()); RpcStatus status = RpcStatus.valueOf(responseMsg.getStatus());
ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap.get(responseMsg.getRequestId()); ToDeviceRpcRequestMetadata md = toDeviceRpcPendingMap.get(responseMsg.getRequestId());
@ -585,6 +581,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
if (status.equals(RpcStatus.DELIVERED)) { if (status.equals(RpcStatus.DELIVERED)) {
if (md.getMsg().getMsg().isOneway()) { if (md.getMsg().getMsg().isOneway()) {
toDeviceRpcPendingMap.remove(responseMsg.getRequestId()); toDeviceRpcPendingMap.remove(responseMsg.getRequestId());
if (rpcSequenceEnabled) {
systemContext.getTbCoreDeviceRpcService().processRpcResponseFromDeviceActor(new FromDeviceRpcResponse(rpcId, null, null));
}
} else { } else {
md.setDelivered(true); md.setDelivered(true);
} }
@ -597,7 +596,9 @@ class DeviceActorMessageProcessor extends AbstractContextAwareMsgProcessor {
} }
} }
systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), status, null); if (md.getMsg().getMsg().isPersisted()) {
systemContext.getTbRpcService().save(tenantId, new RpcId(rpcId), status, null);
}
if (status != RpcStatus.SENT) { if (status != RpcStatus.SENT) {
sendNextPendingRequest(context); sendNextPendingRequest(context);
} }

5
application/src/main/java/org/thingsboard/server/controller/AbstractRpcController.java

@ -75,8 +75,8 @@ public abstract class AbstractRpcController extends BaseController {
SecurityUser currentUser = getCurrentUser(); SecurityUser currentUser = getCurrentUser();
TenantId tenantId = currentUser.getTenantId(); TenantId tenantId = currentUser.getTenantId();
final DeferredResult<ResponseEntity> response = new DeferredResult<>(); final DeferredResult<ResponseEntity> response = new DeferredResult<>();
long timeout = rpcRequestBody.has("timeout") ? rpcRequestBody.get("timeout").asLong() : defaultTimeout; long timeout = rpcRequestBody.has(DataConstants.TIMEOUT) ? rpcRequestBody.get(DataConstants.TIMEOUT).asLong() : defaultTimeout;
long expTime = System.currentTimeMillis() + Math.max(minTimeout, timeout); long expTime = rpcRequestBody.has(DataConstants.EXPIRATION_TIME) ? rpcRequestBody.get(DataConstants.EXPIRATION_TIME).asLong() : System.currentTimeMillis() + Math.max(minTimeout, timeout);
UUID rpcRequestUUID = rpcRequestBody.has("requestUUID") ? UUID.fromString(rpcRequestBody.get("requestUUID").asText()) : UUID.randomUUID(); UUID rpcRequestUUID = rpcRequestBody.has("requestUUID") ? UUID.fromString(rpcRequestBody.get("requestUUID").asText()) : UUID.randomUUID();
boolean persisted = rpcRequestBody.has(DataConstants.PERSISTENT) && rpcRequestBody.get(DataConstants.PERSISTENT).asBoolean(); boolean persisted = rpcRequestBody.has(DataConstants.PERSISTENT) && rpcRequestBody.get(DataConstants.PERSISTENT).asBoolean();
String additionalInfo = JacksonUtil.toString(rpcRequestBody.get(DataConstants.ADDITIONAL_INFO)); String additionalInfo = JacksonUtil.toString(rpcRequestBody.get(DataConstants.ADDITIONAL_INFO));
@ -88,6 +88,7 @@ public abstract class AbstractRpcController extends BaseController {
deviceId, deviceId,
oneWay, oneWay,
expTime, expTime,
timeout,
body, body,
persisted, persisted,
additionalInfo additionalInfo

2
application/src/main/java/org/thingsboard/server/service/rpc/DefaultTbRuleEngineRpcService.java

@ -101,7 +101,7 @@ public class DefaultTbRuleEngineRpcService implements TbRuleEngineDeviceRpcServi
@Override @Override
public void sendRpcRequestToDevice(RuleEngineDeviceRpcRequest src, Consumer<RuleEngineDeviceRpcResponse> consumer) { public void sendRpcRequestToDevice(RuleEngineDeviceRpcRequest src, Consumer<RuleEngineDeviceRpcResponse> consumer) {
ToDeviceRpcRequest request = new ToDeviceRpcRequest(src.getRequestUUID(), src.getTenantId(), src.getDeviceId(), ToDeviceRpcRequest request = new ToDeviceRpcRequest(src.getRequestUUID(), src.getTenantId(), src.getDeviceId(),
src.isOneway(), src.getExpirationTime(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo()); src.isOneway(), src.getExpirationTime(), src.getTimeout(), new ToDeviceRpcRequestBody(src.getMethod(), src.getBody()), src.isPersisted(), src.getAdditionalInfo());
forwardRpcRequestToDeviceActor(request, response -> { forwardRpcRequestToDeviceActor(request, response -> {
if (src.isRestApiCall()) { if (src.isRestApiCall()) {
sendRpcResponseToTbCore(src.getOriginServiceId(), response); sendRpcResponseToTbCore(src.getOriginServiceId(), response);

17
application/src/test/java/org/thingsboard/server/transport/mqtt/rpc/AbstractMqttServerSideRpcIntegrationTest.java

@ -28,15 +28,16 @@ import org.eclipse.paho.client.mqttv3.MqttCallback;
import org.eclipse.paho.client.mqttv3.MqttException; import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.MqttMessage; import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.junit.Assert; import org.junit.Assert;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.TransportPayloadType; import org.thingsboard.server.common.data.TransportPayloadType;
import org.thingsboard.server.common.data.device.profile.MqttTopics; import org.thingsboard.server.common.data.device.profile.MqttTopics;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest; import org.thingsboard.server.transport.mqtt.AbstractMqttIntegrationTest;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Arrays; import java.util.Arrays;
import java.util.List; import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch; import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
@ -120,12 +121,13 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM
} }
MqttAsyncClient client = getMqttAsyncClient(accessToken); MqttAsyncClient client = getMqttAsyncClient(accessToken);
client.setManualAcks(true);
CountDownLatch latch = new CountDownLatch(10); CountDownLatch latch = new CountDownLatch(10);
TestSequenceMqttCallback callback = new TestSequenceMqttCallback(client, latch, result); TestSequenceMqttCallback callback = new TestSequenceMqttCallback(client, latch, result);
client.setCallback(callback); client.setCallback(callback);
client.subscribe(MqttTopics.DEVICE_RPC_REQUESTS_SUB_TOPIC, 1); client.subscribe(MqttTopics.DEVICE_RPC_REQUESTS_SUB_TOPIC, 1);
latch.await(30, TimeUnit.SECONDS); latch.await(10, TimeUnit.SECONDS);
Assert.assertEquals(expected, result); Assert.assertEquals(expected, result);
} }
@ -246,8 +248,7 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM
private final MqttAsyncClient client; private final MqttAsyncClient client;
private final CountDownLatch latch; private final CountDownLatch latch;
private final List<String> expected; private final List<String> expected;
private Integer qoS;
TestSequenceMqttCallback(MqttAsyncClient client, CountDownLatch latch, List<String> expected) { TestSequenceMqttCallback(MqttAsyncClient client, CountDownLatch latch, List<String> expected) {
this.client = client; this.client = client;
@ -255,10 +256,6 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM
this.expected = expected; this.expected = expected;
} }
int getQoS() {
return qoS;
}
@Override @Override
public void connectionLost(Throwable throwable) { public void connectionLost(Throwable throwable) {
} }
@ -268,7 +265,9 @@ public abstract class AbstractMqttServerSideRpcIntegrationTest extends AbstractM
log.info("Message Arrived: " + Arrays.toString(mqttMessage.getPayload())); log.info("Message Arrived: " + Arrays.toString(mqttMessage.getPayload()));
expected.add(new String(mqttMessage.getPayload())); expected.add(new String(mqttMessage.getPayload()));
String responseTopic = requestTopic.replace("request", "response"); String responseTopic = requestTopic.replace("request", "response");
qoS = mqttMessage.getQos(); var qoS = mqttMessage.getQos();
client.messageArrivedComplete(mqttMessage.getId(), qoS);
client.publish(responseTopic, processMessageArrived(requestTopic, mqttMessage)); client.publish(responseTopic, processMessageArrived(requestTopic, mqttMessage));
latch.countDown(); latch.countDown();
} }

5
common/cluster-api/src/main/proto/queue.proto

@ -334,6 +334,7 @@ message ToDeviceRpcRequestMsg {
int64 requestIdLSB = 6; int64 requestIdLSB = 6;
bool oneway = 7; bool oneway = 7;
bool persisted = 8; bool persisted = 8;
int64 timeout = 9;
} }
message ToDeviceRpcResponseMsg { message ToDeviceRpcResponseMsg {
@ -346,7 +347,7 @@ message UplinkNotificationMsg {
int64 uplinkTs = 1; int64 uplinkTs = 1;
} }
message ToDevicePersistedRpcResponseMsg { message ToDeviceRpcResponseStatusMsg {
int32 requestId = 1; int32 requestId = 1;
int64 requestIdMSB = 2; int64 requestIdMSB = 2;
int64 requestIdLSB = 3; int64 requestIdLSB = 3;
@ -456,7 +457,7 @@ message TransportToDeviceActorMsg {
SubscriptionInfoProto subscriptionInfo = 7; SubscriptionInfoProto subscriptionInfo = 7;
ClaimDeviceMsg claimDevice = 8; ClaimDeviceMsg claimDevice = 8;
ProvisionDeviceRequestMsg provisionDevice = 9; ProvisionDeviceRequestMsg provisionDevice = 9;
ToDevicePersistedRpcResponseMsg persistedRpcResponseMsg = 10; ToDeviceRpcResponseStatusMsg rpcResponseStatusMsg = 10;
SendPendingRPCMsg sendPendingRPC = 11; SendPendingRPCMsg sendPendingRPC = 11;
UplinkNotificationMsg uplinkNotificationMsg = 12; UplinkNotificationMsg uplinkNotificationMsg = 12;
} }

2
common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java

@ -36,6 +36,8 @@ public class DataConstants {
public static final String ALARM_CONDITION_REPEATS = "alarmConditionRepeats"; public static final String ALARM_CONDITION_REPEATS = "alarmConditionRepeats";
public static final String ALARM_CONDITION_DURATION = "alarmConditionDuration"; public static final String ALARM_CONDITION_DURATION = "alarmConditionDuration";
public static final String PERSISTENT = "persistent"; public static final String PERSISTENT = "persistent";
public static final String TIMEOUT = "timeout";
public static final String EXPIRATION_TIME = "expirationTime";
public static final String ADDITIONAL_INFO = "additionalInfo"; public static final String ADDITIONAL_INFO = "additionalInfo";
public static final String COAP_TRANSPORT_NAME = "COAP"; public static final String COAP_TRANSPORT_NAME = "COAP";
public static final String LWM2M_TRANSPORT_NAME = "LWM2M"; public static final String LWM2M_TRANSPORT_NAME = "LWM2M";

1
common/message/src/main/java/org/thingsboard/server/common/msg/rpc/ToDeviceRpcRequest.java

@ -34,6 +34,7 @@ public class ToDeviceRpcRequest implements Serializable {
private final DeviceId deviceId; private final DeviceId deviceId;
private final boolean oneway; private final boolean oneway;
private final long expirationTime; private final long expirationTime;
private final long timeout;
private final ToDeviceRpcRequestBody body; private final ToDeviceRpcRequestBody body;
private final boolean persisted; private final boolean persisted;
@JsonIgnore @JsonIgnore

25
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java

@ -525,17 +525,26 @@ public class DefaultCoapClientContext implements CoapClientContext {
Response response = state.getAdaptor().convertToPublish(conRequest, msg, state.getConfiguration().getRpcRequestDynamicMessageBuilder()); Response response = state.getAdaptor().convertToPublish(conRequest, msg, state.getConfiguration().getRpcRequestDynamicMessageBuilder());
int requestId = getNextMsgId(); int requestId = getNextMsgId();
response.setMID(requestId); response.setMID(requestId);
if (msg.getPersisted() && conRequest) { if (conRequest) {
transportContext.getRpcAwaitingAck().put(requestId, msg); transportContext.getRpcAwaitingAck().put(requestId, msg);
transportContext.getScheduler().schedule(() -> { transportContext.getScheduler().schedule(() -> {
transportContext.getRpcAwaitingAck().remove(requestId); TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(requestId);
}, Math.max(0, msg.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); if (rpcRequestMsg != null) {
transportService.process(state.getSession(), msg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY);
}
}, Math.max(0, Math.min(msg.getTimeout(), msg.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS);
response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> {
TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id);
if (rpcRequestMsg != null) { if (rpcRequestMsg != null) {
transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); transportService.process(state.getSession(), rpcRequestMsg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY);
} }
}, null)); }, id -> {
TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id);
if (rpcRequestMsg != null) {
transportService.process(state.getSession(), msg, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY);
}
}));
} }
if (conRequest) { if (conRequest) {
response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> awake(state), id -> asleep(state))); response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> awake(state), id -> asleep(state)));
@ -554,11 +563,11 @@ public class DefaultCoapClientContext implements CoapClientContext {
transportService.process(state.getSession(), transportService.process(state.getSession(),
TransportProtos.ToDeviceRpcResponseMsg.newBuilder() TransportProtos.ToDeviceRpcResponseMsg.newBuilder()
.setRequestId(msg.getRequestId()).setError(error).build(), TransportServiceCallback.EMPTY); .setRequestId(msg.getRequestId()).setError(error).build(), TransportServiceCallback.EMPTY);
} else if (msg.getPersisted() && sent) { } else if (sent) {
if (conRequest) { if (!conRequest) {
transportService.process(state.getSession(), msg, RpcStatus.SENT, TransportServiceCallback.EMPTY);
} else {
transportService.process(state.getSession(), msg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY); transportService.process(state.getSession(), msg, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY);
} else if (msg.getPersisted()) {
transportService.process(state.getSession(), msg, RpcStatus.SENT, TransportServiceCallback.EMPTY);
} }
} }
} }

29
common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

@ -850,24 +850,27 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement
try { try {
deviceSessionCtx.getPayloadAdaptor().convertToPublish(deviceSessionCtx, rpcRequest).ifPresent(payload -> { deviceSessionCtx.getPayloadAdaptor().convertToPublish(deviceSessionCtx, rpcRequest).ifPresent(payload -> {
int msgId = ((MqttPublishMessage) payload).variableHeader().packetId(); int msgId = ((MqttPublishMessage) payload).variableHeader().packetId();
if (rpcRequest.getPersisted() && isAckExpected(payload)) { if (isAckExpected(payload)) {
rpcAwaitingAck.put(msgId, rpcRequest); rpcAwaitingAck.put(msgId, rpcRequest);
context.getScheduler().schedule(() -> { context.getScheduler().schedule(() -> {
rpcAwaitingAck.remove(msgId); TransportProtos.ToDeviceRpcRequestMsg msg = rpcAwaitingAck.remove(msgId);
}, Math.max(0, rpcRequest.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); if (msg != null) {
transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.TIMEOUT, TransportServiceCallback.EMPTY);
}
}, Math.max(0, Math.min(rpcRequest.getTimeout(), rpcRequest.getExpirationTime() - System.currentTimeMillis())), TimeUnit.MILLISECONDS);
} }
var cf = publish(payload, deviceSessionCtx); var cf = publish(payload, deviceSessionCtx);
if (rpcRequest.getPersisted()) { cf.addListener(result -> {
cf.addListener(result -> { if (result.cause() == null) {
if (result.cause() == null) { if (!isAckExpected(payload)) {
if (isAckExpected(payload)) { transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY);
transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY); } else if (rpcRequest.getPersisted()) {
} else { transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.SENT, TransportServiceCallback.EMPTY);
transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, RpcStatus.DELIVERED, TransportServiceCallback.EMPTY);
}
} }
}); } else {
} // TODO: send error
}
});
}); });
} catch (Exception e) { } catch (Exception e) {
transportService.process(deviceSessionCtx.getSessionInfo(), transportService.process(deviceSessionCtx.getSessionInfo(),

24
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/service/DefaultTransportService.java

@ -581,19 +581,17 @@ public class DefaultTransportService implements TransportService {
@Override @Override
public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, RpcStatus rpcStatus, TransportServiceCallback<Void> callback) { public void process(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToDeviceRpcRequestMsg msg, RpcStatus rpcStatus, TransportServiceCallback<Void> callback) {
if (msg.getPersisted()) { TransportProtos.ToDeviceRpcResponseStatusMsg responseMsg = TransportProtos.ToDeviceRpcResponseStatusMsg.newBuilder()
TransportProtos.ToDevicePersistedRpcResponseMsg responseMsg = TransportProtos.ToDevicePersistedRpcResponseMsg.newBuilder() .setRequestId(msg.getRequestId())
.setRequestId(msg.getRequestId()) .setRequestIdLSB(msg.getRequestIdLSB())
.setRequestIdLSB(msg.getRequestIdLSB()) .setRequestIdMSB(msg.getRequestIdMSB())
.setRequestIdMSB(msg.getRequestIdMSB()) .setStatus(rpcStatus.name())
.setStatus(rpcStatus.name()) .build();
.build();
if (checkLimits(sessionInfo, responseMsg, callback)) {
if (checkLimits(sessionInfo, responseMsg, callback)) { reportActivityInternal(sessionInfo);
reportActivityInternal(sessionInfo); sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setRpcResponseStatusMsg(responseMsg).build(),
sendToDeviceActor(sessionInfo, TransportToDeviceActorMsg.newBuilder().setSessionInfo(sessionInfo).setPersistedRpcResponseMsg(responseMsg).build(), new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY));
new ApiStatsProxyCallback<>(getTenantId(sessionInfo), getCustomerId(sessionInfo), 1, TransportServiceCallback.EMPTY));
}
} }
} }

1
rule-engine/rule-engine-api/src/main/java/org/thingsboard/rule/engine/api/RuleEngineDeviceRpcRequest.java

@ -39,6 +39,7 @@ public final class RuleEngineDeviceRpcRequest {
private final String method; private final String method;
private final String body; private final String body;
private final long expirationTime; private final long expirationTime;
private final long timeout;
private final boolean restApiCall; private final boolean restApiCall;
private final String additionalInfo; private final String additionalInfo;

6
rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/rpc/TbSendRPCRequestNode.java

@ -89,9 +89,12 @@ public class TbSendRPCRequestNode implements TbNode {
tmp = msg.getMetaData().getValue("originServiceId"); tmp = msg.getMetaData().getValue("originServiceId");
String originServiceId = !StringUtils.isEmpty(tmp) ? tmp : null; String originServiceId = !StringUtils.isEmpty(tmp) ? tmp : null;
tmp = msg.getMetaData().getValue("expirationTime"); tmp = msg.getMetaData().getValue(DataConstants.EXPIRATION_TIME);
long expirationTime = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : (System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds())); long expirationTime = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : (System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds()));
tmp = msg.getMetaData().getValue(DataConstants.TIMEOUT);
long timeout = !StringUtils.isEmpty(tmp) ? Long.parseLong(tmp) : TimeUnit.SECONDS.toMillis(config.getTimeoutInSeconds());
String params; String params;
JsonElement paramsEl = json.get("params"); JsonElement paramsEl = json.get("params");
if (paramsEl.isJsonPrimitive()) { if (paramsEl.isJsonPrimitive()) {
@ -112,6 +115,7 @@ public class TbSendRPCRequestNode implements TbNode {
.requestUUID(requestUUID) .requestUUID(requestUUID)
.originServiceId(originServiceId) .originServiceId(originServiceId)
.expirationTime(expirationTime) .expirationTime(expirationTime)
.timeout(timeout)
.restApiCall(restApiCall) .restApiCall(restApiCall)
.persisted(persisted) .persisted(persisted)
.additionalInfo(additionalInfo) .additionalInfo(additionalInfo)

Loading…
Cancel
Save