Browse Source

PR Review + minor logging

pull/4910/head
Andrii Shvaika 5 years ago
parent
commit
e914425b22
  1. 4
      application/src/main/java/org/thingsboard/server/service/ota/DefaultOtaPackageStateService.java
  2. 20
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
  3. 8
      common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/MqttTransportHandler.java

4
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);
}
});

20
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);

8
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);
});

Loading…
Cancel
Save