Browse Source

CoAP ack/confirmable mess fixing

pull/5503/head
Andrii Shvaika 5 years ago
committed by Andrew Shvayka
parent
commit
de1a56a488
  1. 4
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java
  2. 8
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java
  3. 18
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java
  4. 18
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/ProtoCoapAdaptor.java
  5. 4
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/CoapOkCallback.java
  6. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/GetAttributesSyncSessionCallback.java
  7. 2
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/ToServerRpcSyncSessionCallback.java
  8. 9
      common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/client/DefaultCoapClientContext.java

4
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/CoapTransportResource.java

@ -223,9 +223,7 @@ public class CoapTransportResource extends AbstractCoapTransportResource {
return;
}
transportService.process(DeviceTransportType.COAP, TransportProtos.ValidateDeviceTokenRequestMsg.newBuilder().setToken(credentials.get().getCredentialsId()).build(),
new CoapDeviceAuthCallback(exchange, (deviceCredentials, deviceProfile) -> {
processRequest(exchange, type, request, deviceCredentials, deviceProfile);
}));
new CoapDeviceAuthCallback(exchange, (deviceCredentials, deviceProfile) -> processRequest(exchange, type, request, deviceCredentials, deviceProfile)));
}
private void processRequest(CoapExchange exchange, SessionMsgType type, Request request, ValidateDeviceCredentialsResponse deviceCredentials, DeviceProfile deviceProfile) {

8
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/CoapTransportAdaptor.java

@ -39,13 +39,13 @@ public interface CoapTransportAdaptor {
TransportProtos.ClaimDeviceMsg convertToClaimDevice(UUID sessionId, Request inbound, TransportProtos.SessionInfoProto sessionInfo) throws AdaptorException;
Response convertToPublish(boolean isConfirmable, TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException;
Response convertToPublish(TransportProtos.GetAttributeResponseMsg responseMsg) throws AdaptorException;
Response convertToPublish(boolean isConfirmable, TransportProtos.AttributeUpdateNotificationMsg notificationMsg) throws AdaptorException;
Response convertToPublish(TransportProtos.AttributeUpdateNotificationMsg notificationMsg) throws AdaptorException;
Response convertToPublish(boolean isConfirmable, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) throws AdaptorException;
Response convertToPublish(TransportProtos.ToDeviceRpcRequestMsg rpcRequest, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) throws AdaptorException;
Response convertToPublish(boolean isConfirmable, TransportProtos.ToServerRpcResponseMsg msg) throws AdaptorException;
Response convertToPublish(TransportProtos.ToServerRpcResponseMsg msg) throws AdaptorException;
ProvisionDeviceRequestMsg convertToProvisionRequestMsg(UUID sessionId, Request inbound) throws AdaptorException;

18
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/JsonCoapAdaptor.java

@ -93,21 +93,20 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor {
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.AttributeUpdateNotificationMsg msg) throws AdaptorException {
return getObserveNotification(isConfirmable, JsonConverter.toJson(msg));
public Response convertToPublish(TransportProtos.AttributeUpdateNotificationMsg msg) throws AdaptorException {
return getObserveNotification(JsonConverter.toJson(msg));
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.ToDeviceRpcRequestMsg msg, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) throws AdaptorException {
return getObserveNotification(isConfirmable, JsonConverter.toJson(msg, true));
public Response convertToPublish(TransportProtos.ToDeviceRpcRequestMsg msg, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) throws AdaptorException {
return getObserveNotification(JsonConverter.toJson(msg, true));
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.ToServerRpcResponseMsg msg) throws AdaptorException {
public Response convertToPublish(TransportProtos.ToServerRpcResponseMsg msg) throws AdaptorException {
Response response = new Response(CoAP.ResponseCode.CONTENT);
JsonElement result = JsonConverter.toJson(msg);
response.setPayload(result.toString());
response.setConfirmable(isConfirmable);
return response;
}
@ -122,11 +121,10 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor {
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException {
public Response convertToPublish(TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException {
if (msg.getSharedStateMsg()) {
if (StringUtils.isEmpty(msg.getError())) {
Response response = new Response(CoAP.ResponseCode.CONTENT);
response.setConfirmable(isConfirmable);
TransportProtos.AttributeUpdateNotificationMsg notificationMsg = TransportProtos.AttributeUpdateNotificationMsg.newBuilder().addAllSharedUpdated(msg.getSharedAttributeListList()).build();
JsonObject result = JsonConverter.toJson(notificationMsg);
response.setPayload(result.toString());
@ -139,7 +137,6 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor {
return new Response(CoAP.ResponseCode.NOT_FOUND);
} else {
Response response = new Response(CoAP.ResponseCode.CONTENT);
response.setConfirmable(isConfirmable);
JsonObject result = JsonConverter.toJson(msg);
response.setPayload(result.toString());
return response;
@ -147,10 +144,9 @@ public class JsonCoapAdaptor implements CoapTransportAdaptor {
}
}
private Response getObserveNotification(boolean confirmable, JsonElement json) {
private Response getObserveNotification(JsonElement json) {
Response response = new Response(CoAP.ResponseCode.CONTENT);
response.setPayload(json.toString());
response.setConfirmable(confirmable);
return response;
}

18
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/adaptors/ProtoCoapAdaptor.java

@ -113,29 +113,27 @@ public class ProtoCoapAdaptor implements CoapTransportAdaptor {
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.AttributeUpdateNotificationMsg msg) throws AdaptorException {
return getObserveNotification(isConfirmable, msg.toByteArray());
public Response convertToPublish(TransportProtos.AttributeUpdateNotificationMsg msg) throws AdaptorException {
return getObserveNotification(msg.toByteArray());
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.ToDeviceRpcRequestMsg rpcRequest, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) throws AdaptorException {
return getObserveNotification(isConfirmable, ProtoConverter.convertToRpcRequest(rpcRequest, rpcRequestDynamicMessageBuilder));
public Response convertToPublish(TransportProtos.ToDeviceRpcRequestMsg rpcRequest, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) throws AdaptorException {
return getObserveNotification(ProtoConverter.convertToRpcRequest(rpcRequest, rpcRequestDynamicMessageBuilder));
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.ToServerRpcResponseMsg msg) throws AdaptorException {
public Response convertToPublish(TransportProtos.ToServerRpcResponseMsg msg) throws AdaptorException {
Response response = new Response(CoAP.ResponseCode.CONTENT);
response.setConfirmable(isConfirmable);
response.setPayload(msg.toByteArray());
return response;
}
@Override
public Response convertToPublish(boolean isConfirmable, TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException {
public Response convertToPublish(TransportProtos.GetAttributeResponseMsg msg) throws AdaptorException {
if (msg.getSharedStateMsg()) {
if (StringUtils.isEmpty(msg.getError())) {
Response response = new Response(CoAP.ResponseCode.CONTENT);
response.setConfirmable(isConfirmable);
TransportProtos.AttributeUpdateNotificationMsg notificationMsg = TransportProtos.AttributeUpdateNotificationMsg.newBuilder().addAllSharedUpdated(msg.getSharedAttributeListList()).build();
response.setPayload(notificationMsg.toByteArray());
return response;
@ -147,17 +145,15 @@ public class ProtoCoapAdaptor implements CoapTransportAdaptor {
return new Response(CoAP.ResponseCode.NOT_FOUND);
} else {
Response response = new Response(CoAP.ResponseCode.CONTENT);
response.setConfirmable(isConfirmable);
response.setPayload(msg.toByteArray());
return response;
}
}
}
private Response getObserveNotification(boolean confirmable, byte[] notification) {
private Response getObserveNotification(byte[] notification) {
Response response = new Response(CoAP.ResponseCode.CONTENT);
response.setPayload(notification);
response.setConfirmable(confirmable);
return response;
}

4
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/CoapOkCallback.java

@ -34,9 +34,7 @@ public class CoapOkCallback implements TransportServiceCallback<Void> {
@Override
public void onSuccess(Void msg) {
Response response = new Response(onSuccessResponse);
response.setConfirmable(isConRequest());
exchange.respond(response);
exchange.respond(new Response(onSuccessResponse));
}
@Override

2
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/GetAttributesSyncSessionCallback.java

@ -34,7 +34,7 @@ public class GetAttributesSyncSessionCallback extends AbstractSyncSessionCallbac
@Override
public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg msg) {
try {
respond(state.getAdaptor().convertToPublish(request.isConfirmable(), msg));
respond(state.getAdaptor().convertToPublish(msg));
} catch (AdaptorException e) {
log.trace("[{}] Failed to reply due to error", state.getDeviceId(), e);
exchange.respond(new Response(CoAP.ResponseCode.INTERNAL_SERVER_ERROR));

2
common/transport/coap/src/main/java/org/thingsboard/server/transport/coap/callback/ToServerRpcSyncSessionCallback.java

@ -33,7 +33,7 @@ public class ToServerRpcSyncSessionCallback extends AbstractSyncSessionCallback
@Override
public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) {
try {
respond(state.getAdaptor().convertToPublish(request.isConfirmable(), toServerResponse));
respond(state.getAdaptor().convertToPublish(toServerResponse));
} catch (AdaptorException e) {
log.trace("Failed to reply due to error", e);
exchange.respond(CoAP.ResponseCode.INTERNAL_SERVER_ERROR);

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

@ -435,8 +435,7 @@ public class DefaultCoapClientContext implements CoapClientContext {
TbCoapObservationState attrs = state.getAttrs();
if (attrs != null) {
try {
boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getAttrs());
Response response = state.getAdaptor().convertToPublish(conRequest, msg);
Response response = state.getAdaptor().convertToPublish(msg);
respond(attrs.getExchange(), response, state.getContentFormat());
} catch (AdaptorException e) {
log.trace("Failed to reply due to error", e);
@ -466,7 +465,8 @@ public class DefaultCoapClientContext implements CoapClientContext {
try {
boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getAttrs());
int requestId = getNextMsgId();
Response response = state.getAdaptor().convertToPublish(conRequest, msg);
Response response = state.getAdaptor().convertToPublish(msg);
response.setConfirmable(conRequest);
response.setMID(requestId);
if (conRequest) {
response.addMessageObserver(new TbCoapMessageObserver(requestId, id -> awake(state), id -> asleep(state)));
@ -527,7 +527,8 @@ public class DefaultCoapClientContext implements CoapClientContext {
String error = null;
boolean conRequest = AbstractSyncSessionCallback.isConRequest(state.getRpc());
try {
Response response = state.getAdaptor().convertToPublish(conRequest, msg, state.getConfiguration().getRpcRequestDynamicMessageBuilder());
Response response = state.getAdaptor().convertToPublish(msg, state.getConfiguration().getRpcRequestDynamicMessageBuilder());
response.setConfirmable(conRequest);
int requestId = getNextMsgId();
response.setMID(requestId);
if (conRequest) {

Loading…
Cancel
Save