diff --git a/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java index 76797137b7..2d6de7cfdf 100644 --- a/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java @@ -273,13 +273,13 @@ public class DefaultOtaPackageStateService implements OtaPackageStateService { telemetryService.saveAndNotify(tenantId, deviceId, Collections.singletonList(status), new FutureCallback<>() { @Override public void onSuccess(@Nullable Void tmp) { - log.trace("[{}] Success save telemetry with target firmware for device!", deviceId); + log.trace("[{}] Success save telemetry with target {} for device!", deviceId, otaPackage); updateAttributes(device, otaPackage, ts, tenantId, deviceId, otaPackageType); } @Override public void onFailure(Throwable t) { - log.error("[{}] Failed to save telemetry with target firmware for device!", deviceId, t); + log.error("[{}] Failed to save telemetry with target {} for device!", deviceId, otaPackage, t); updateAttributes(device, otaPackage, ts, tenantId, deviceId, otaPackageType); } }); diff --git a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java index 33834eed22..2facdcd305 100644 --- a/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java +++ b/common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java @@ -497,8 +497,8 @@ public class CoapTransportResource extends AbstractCoapTransportResource { int requestId = getNextMsgId(); response.setMID(requestId); - if (isConRequest()) { - if (msg.getPersisted()) { + if (msg.getPersisted()) { + if (isConRequest()) { transportContext.getRpcAwaitingAck().put(requestId, msg); transportContext.getScheduler().schedule(() -> { TransportProtos.ToDeviceRpcRequestMsg awaitingAckMsg = transportContext.getRpcAwaitingAck().remove(requestId); @@ -506,15 +506,15 @@ public class CoapTransportResource extends AbstractCoapTransportResource { transportService.process(sessionInfo, msg, true, TransportServiceCallback.EMPTY); } }, Math.max(0, msg.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); + response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { + TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); + if (rpcRequestMsg != null) { + transportService.process(sessionInfo, rpcRequestMsg, false, TransportServiceCallback.EMPTY); + } + })); + } else { + transportService.process(sessionInfo, msg, false, TransportServiceCallback.EMPTY); } - response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> { - TransportProtos.ToDeviceRpcRequestMsg rpcRequestMsg = transportContext.getRpcAwaitingAck().remove(id); - if (rpcRequestMsg != null) { - transportService.process(sessionInfo, rpcRequestMsg, false, TransportServiceCallback.EMPTY); - } - })); - } else if (msg.getPersisted()) { - transportService.process(sessionInfo, msg, false, TransportServiceCallback.EMPTY); } exchange.respond(response); diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java index c42bb0163b..e09a47972b 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java @@ -825,8 +825,8 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement try { deviceSessionCtx.getPayloadAdaptor().convertToPublish(deviceSessionCtx, rpcRequest).ifPresent(payload -> { int msgId = ((MqttPublishMessage) payload).variableHeader().packetId(); - if (isAckExpected(payload)) { - if (rpcRequest.getPersisted()) { + if (rpcRequest.getPersisted()) { + if (isAckExpected(payload)) { rpcAwaitingAck.put(msgId, rpcRequest); context.getScheduler().schedule(() -> { TransportProtos.ToDeviceRpcRequestMsg awaitingAckMsg = rpcAwaitingAck.remove(msgId); @@ -834,9 +834,9 @@ public class MqttTransportHandler extends ChannelInboundHandlerAdapter implement transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, true, TransportServiceCallback.EMPTY); } }, Math.max(0, rpcRequest.getExpirationTime() - System.currentTimeMillis()), TimeUnit.MILLISECONDS); + } else { + transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, false, TransportServiceCallback.EMPTY); } - } else if (rpcRequest.getPersisted()) { - transportService.process(deviceSessionCtx.getSessionInfo(), rpcRequest, false, TransportServiceCallback.EMPTY); } publish(payload, deviceSessionCtx); });