@ -49,6 +49,7 @@ import org.eclipse.leshan.server.registration.Registration;
import org.springframework.stereotype.Service ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.server.common.data.device.data.lwm2m.ObjectAttributes ;
import org.thingsboard.server.common.transport.util.JsonUtils ;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent ;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig ;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext ;
@ -88,7 +89,6 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
private final LwM2mTransportContext context ;
private final LwM2MTransportServerConfig config ;
private final LwM2mClientContext lwM2mClientContext ;
@PostConstruct
public void init ( ) {
@ -98,14 +98,14 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
}
@Override
public void sendReadRequest ( LwM2mClient client , TbLwM2MReadRequest request , DownlinkRequestCallback < ReadResponse > callback ) {
public void sendReadRequest ( LwM2mClient client , TbLwM2MReadRequest request , DownlinkRequestCallback < ReadRequest , ReadRe sponse > callback ) {
validateVersionedId ( client , request ) ;
ReadRequest downlink = new ReadRequest ( getContentFormat ( client , request ) , request . getObjectId ( ) ) ;
sendRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
}
@Override
public void sendObserveRequest ( LwM2mClient client , TbLwM2MObserveRequest request , DownlinkRequestCallback < ObserveResponse > callback ) {
public void sendObserveRequest ( LwM2mClient client , TbLwM2MObserveRequest request , DownlinkRequestCallback < ObserveRequest , ObserveRe sponse > callback ) {
validateVersionedId ( client , request ) ;
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
Set < Observation > observations = context . getServer ( ) . getObservationService ( ) . getObservations ( client . getRegistration ( ) ) ;
@ -127,19 +127,19 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
}
@Override
public void sendObserveAllRequest ( LwM2mClient client , TbLwM2MObserveAllRequest request , DownlinkRequestCallback < Set < String > > callback ) {
public void sendObserveAllRequest ( LwM2mClient client , TbLwM2MObserveAllRequest request , DownlinkRequestCallback < TbLwM2MObserveAllRequest , Set < String > > callback ) {
Set < Observation > observations = context . getServer ( ) . getObservationService ( ) . getObservations ( client . getRegistration ( ) ) ;
Set < String > paths = observations . stream ( ) . map ( observation - > observation . getPath ( ) . toString ( ) ) . collect ( Collectors . toUnmodifiableSet ( ) ) ;
callback . onSuccess ( paths ) ;
callback . onSuccess ( request , paths ) ;
}
@Override
public void sendDiscoverAllRequest ( LwM2mClient client , TbLwM2MDiscoverAllRequest request , DownlinkRequestCallback < List < Link > > callback ) {
callback . onSuccess ( Arrays . asList ( client . getRegistration ( ) . getSortedObjectLinks ( ) ) ) ;
public void sendDiscoverAllRequest ( LwM2mClient client , TbLwM2MDiscoverAllRequest request , DownlinkRequestCallback < TbLwM2MDiscoverAllRequest , List < Link > > callback ) {
callback . onSuccess ( request , Arrays . asList ( client . getRegistration ( ) . getSortedObjectLinks ( ) ) ) ;
}
@Override
public void sendExecuteRequest ( LwM2mClient client , TbLwM2MExecuteRequest request , DownlinkRequestCallback < ExecuteResponse > callback ) {
public void sendExecuteRequest ( LwM2mClient client , TbLwM2MExecuteRequest request , DownlinkRequestCallback < ExecuteRequest , ExecuteRe sponse > callback ) {
ResourceModel resourceModelExecute = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
if ( resourceModelExecute ! = null ) {
ExecuteRequest downlink ;
@ -153,30 +153,30 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
}
@Override
public void sendDeleteRequest ( LwM2mClient client , TbLwM2MDeleteRequest request , DownlinkRequestCallback < DeleteResponse > callback ) {
public void sendDeleteRequest ( LwM2mClient client , TbLwM2MDeleteRequest request , DownlinkRequestCallback < DeleteRequest , DeleteRe sponse > callback ) {
sendRequest ( client , new DeleteRequest ( request . getObjectId ( ) ) , request . getTimeout ( ) , callback ) ;
}
@Override
public void sendCancelObserveRequest ( LwM2mClient client , TbLwM2MCancelObserveRequest request , DownlinkRequestCallback < Integer > callback ) {
public void sendCancelObserveRequest ( LwM2mClient client , TbLwM2MCancelObserveRequest request , DownlinkRequestCallback < TbLwM2MCancelObserveRequest , Integer > callback ) {
int observeCancelCnt = context . getServer ( ) . getObservationService ( ) . cancelObservations ( client . getRegistration ( ) , request . getObjectId ( ) ) ;
callback . onSuccess ( observeCancelCnt ) ;
callback . onSuccess ( request , observeCancelCnt ) ;
}
@Override
public void sendCancelAllRequest ( LwM2mClient client , TbLwM2MCancelAllRequest request , DownlinkRequestCallback < Integer > callback ) {
public void sendCancelAllRequest ( LwM2mClient client , TbLwM2MCancelAllRequest request , DownlinkRequestCallback < TbLwM2MCancelAllRequest , Integer > callback ) {
int observeCancelCnt = context . getServer ( ) . getObservationService ( ) . cancelObservations ( client . getRegistration ( ) ) ;
callback . onSuccess ( observeCancelCnt ) ;
callback . onSuccess ( request , observeCancelCnt ) ;
}
@Override
public void sendDiscoverRequest ( LwM2mClient client , TbLwM2MDiscoverRequest request , DownlinkRequestCallback < DiscoverResponse > callback ) {
public void sendDiscoverRequest ( LwM2mClient client , TbLwM2MDiscoverRequest request , DownlinkRequestCallback < DiscoverRequest , DiscoverRe sponse > callback ) {
validateVersionedId ( client , request ) ;
sendRequest ( client , new DiscoverRequest ( request . getObjectId ( ) ) , request . getTimeout ( ) , callback ) ;
}
@Override
public void sendWriteAttributesRequest ( LwM2mClient client , TbLwM2MWriteAttributesRequest request , DownlinkRequestCallback < WriteAttributesResponse > callback ) {
public void sendWriteAttributesRequest ( LwM2mClient client , TbLwM2MWriteAttributesRequest request , DownlinkRequestCallback < WriteAttributesRequest , WriteAttributesRe sponse > callback ) {
validateVersionedId ( client , request ) ;
if ( request . getAttributes ( ) = = null ) {
throw new IllegalArgumentException ( "Attributes to write are not specified!" ) ;
@ -196,7 +196,7 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
}
@Override
public void sendWriteReplaceRequest ( LwM2mClient client , TbLwM2MWriteReplaceRequest request , DownlinkRequestCallback < WriteResponse > callback ) {
public void sendWriteReplaceRequest ( LwM2mClient client , TbLwM2MWriteReplaceRequest request , DownlinkRequestCallback < WriteRequest , WriteRe sponse > callback ) {
ResourceModel resourceModelWrite = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
if ( resourceModelWrite ! = null ) {
ContentFormat contentFormat = convertResourceModelTypeToContentFormat ( client , resourceModelWrite . type ) ;
@ -206,15 +206,15 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
path . getObjectId ( ) , path . getObjectInstanceId ( ) , path . getResourceId ( ) , request . getValue ( ) ) ;
sendRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} catch ( Exception e ) {
callback . onError ( e ) ;
callback . onError ( JacksonUtil . toString ( request ) , e ) ;
}
} else {
//TODO: log validation error using callback.
callback . onValidationError ( JacksonUtil . toString ( request ) , "Resource " + request . getVersionedId ( ) + " is not configured in the device profile!" ) ;
}
}
@Override
public void sendWriteUpdateRequest ( LwM2mClient client , TbLwM2MWriteUpdateRequest request , DownlinkRequestCallback < WriteResponse > callback ) {
public void sendWriteUpdateRequest ( LwM2mClient client , TbLwM2MWriteUpdateRequest request , DownlinkRequestCallback < WriteRequest , WriteRe sponse > callback ) {
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
if ( resultIds . isResource ( ) ) {
/ *
@ -238,186 +238,27 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
WriteRequest downlink = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resources ) ;
sendRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} else {
callback . onValidationError ( "No resources to update!" ) ;
callback . onValidationError ( JacksonUtil . toString ( request ) , "No resources to update!" ) ;
}
} else {
callback . onValidationError ( "Update of the root level object is not supported yet!" ) ;
callback . onValidationError ( JacksonUtil . toString ( request ) , "Update of the root level object is not supported yet!" ) ;
}
}
// public void sendAllRequest(LwM2mClient client, String targetIdVer, LwM2mTypeOper typeOper,
// ContentFormat contentFormat, Object params, long timeoutInMs, LwM2mClientRpcRequest lwm2mClientRpcRequest) {
// Registration registration = client.getRegistration();
// try {
// String target = fromVersionedIdToObjectId(targetIdVer);
// if (contentFormat == null) {
// contentFormat = client.getDefaultContentFormat();
// }
// LwM2mPath resultIds = target != null ? new LwM2mPath(target) : null;
// if (!OBSERVE_CANCEL.name().equals(typeOper.name()) && resultIds != null && registration != null && resultIds.getObjectId() >= 0) {
// if (client.isValidObjectVersion(targetIdVer)) {
// timeoutInMs = timeoutInMs > 0 ? timeoutInMs : DEFAULT_TIMEOUT;
// SimpleDownlinkRequest request = createRequest(registration, client, typeOper, contentFormat, target,
// targetIdVer, resultIds, params, lwm2mClientRpcRequest);
// if (request != null) {
// try {
// this.sendRequest(client, request, timeoutInMs, lwm2mClientRpcRequest);
// } catch (ClientSleepingException e) {
// SimpleDownlinkRequest finalRequest = request;
// long finalTimeoutInMs = timeoutInMs;
// LwM2mClientRpcRequest finalRpcRequest = lwm2mClientRpcRequest;
// client.getQueuedRequests().add(() -> sendRequest(client, finalRequest, finalTimeoutInMs, finalRpcRequest));
// } catch (Exception e) {
// log.error("[{}] [{}] [{}] Failed to send downlink.", registration.getEndpoint(), targetIdVer, typeOper.name(), e);
// }
// } else if (WRITE_UPDATE.name().equals(typeOper.name())) {
// if (lwm2mClientRpcRequest != null) {
// String errorMsg = String.format("Path %s params is not valid", targetIdVer);
// handler.sentRpcResponse(lwm2mClientRpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR);
// }
// } else if (WRITE_REPLACE.name().equals(typeOper.name()) || EXECUTE.name().equals(typeOper.name())) {
// if (lwm2mClientRpcRequest != null) {
// String errorMsg = String.format("Path %s object model is absent", targetIdVer);
// handler.sentRpcResponse(lwm2mClientRpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR);
// }
// } else if (!OBSERVE_CANCEL.name().equals(typeOper.name())) {
// log.error("[{}], [{}] - [{}] error SendRequest", registration.getEndpoint(), typeOper.name(), targetIdVer);
// if (lwm2mClientRpcRequest != null) {
// ResourceModel resourceModel = client.getResourceModel(targetIdVer, this.config.getModelProvider());
// String errorMsg = resourceModel == null ? String.format("Path %s not found in object version", targetIdVer) : "SendRequest - null";
// handler.sentRpcResponse(lwm2mClientRpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR);
// }
// }
// } else if (lwm2mClientRpcRequest != null) {
// String errorMsg = String.format("Path %s not found in object version", targetIdVer);
// handler.sentRpcResponse(lwm2mClientRpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR);
// }
// } else {
// switch (typeOper) {
// case OBSERVE_READ_ALL:
// case DISCOVER_ALL:
// Set<String> paths;
// if (OBSERVE_READ_ALL.name().equals(typeOper.name())) {
// Set<Observation> observations = context.getServer().getObservationService().getObservations(registration);
// paths = observations.stream().map(observation -> observation.getPath().toString()).collect(Collectors.toUnmodifiableSet());
// } else {
// assert registration != null;
// Link[] objectLinks = registration.getSortedObjectLinks();
// paths = Arrays.stream(objectLinks).map(Link::toString).collect(Collectors.toUnmodifiableSet());
// }
// String msg = String.format("%s: type operation %s paths - %s", LOG_LW2M_INFO,
// typeOper.name(), paths);
// this.handler.sendLogsToThingsboard(client, msg);
// if (lwm2mClientRpcRequest != null) {
// String valueMsg = String.format("Paths - %s", paths);
// handler.sentRpcResponse(lwm2mClientRpcRequest, CONTENT.name(), valueMsg, LOG_LW2M_VALUE);
// }
// break;
// case OBSERVE_CANCEL:
// case OBSERVE_CANCEL_ALL:
// int observeCancelCnt = 0;
// String observeCancelMsg = null;
// if (OBSERVE_CANCEL.name().equals(typeOper)) {
// observeCancelCnt = context.getServer().getObservationService().cancelObservations(registration, target);
// observeCancelMsg = String.format("%s: type operation %s paths: %s count: %d", LOG_LW2M_INFO,
// OBSERVE_CANCEL.name(), target, observeCancelCnt);
// } else {
// observeCancelCnt = context.getServer().getObservationService().cancelObservations(registration);
// observeCancelMsg = String.format("%s: type operation %s paths: All count: %d", LOG_LW2M_INFO,
// OBSERVE_CANCEL.name(), observeCancelCnt);
// }
// this.afterObserveCancel(client, observeCancelCnt, observeCancelMsg, lwm2mClientRpcRequest);
// break;
// // lwm2mClientRpcRequest != null
// case FW_UPDATE:
// handler.getInfoFirmwareUpdate(client, lwm2mClientRpcRequest);
// break;
// }
// }
// } catch (Exception e) {
// String msg = String.format("%s: type operation %s %s", LOG_LW2M_ERROR,
// typeOper.name(), e.getMessage());
// handler.sendLogsToThingsboard(client, msg);
// if (lwm2mClientRpcRequest != null) {
// String errorMsg = String.format("Path %s type operation %s %s", targetIdVer, typeOper.name(), e.getMessage());
// handler.sentRpcResponse(lwm2mClientRpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR);
// }
// }
// }
private < T extends LwM2mResponse > void sendRequest ( LwM2mClient client , SimpleDownlinkRequest < T > request , long timeoutInMs , DownlinkRequestCallback < T > callback ) {
private < R extends SimpleDownlinkRequest < T > , T extends LwM2mResponse > void sendRequest ( LwM2mClient client , R request , long timeoutInMs , DownlinkRequestCallback < R , T > callback ) {
Registration registration = client . getRegistration ( ) ;
context . getServer ( ) . send ( registration , request , timeoutInMs , response - > {
// if (!client.isInit()) {
// client.initReadValue(this.handler, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration));
// }
responseRequestExecutor . submit ( ( ) - > {
try {
callback . onSuccess ( response ) ;
callback . onSuccess ( request , response ) ;
} catch ( Exception e ) {
log . error ( "[{}] failed to process successful response [{}] " , registration . getEndpoint ( ) , response , e ) ;
}
} ) ;
// if (CoAP.ResponseCode.isSuccess(((Response) response.getCoapResponse()).getCode())) {
// this.handleResponse(client, request.getPath().toString(), response, request, rpcRequest);
// } else {
// String msg = String.format("%s: SendRequest %s: CoapCode - %s Lwm2m code - %d name - %s Resource path - %s", LOG_LW2M_ERROR, request.getClass().getName().toString(),
// ((Response) response.getCoapResponse()).getCode(), response.getCode().getCode(), response.getCode().getName(), request.getPath().toString());
// handler.sendLogsToThingsboard(client, msg);
// log.error("[{}] [{}], [{}] - [{}] [{}] error SendRequest", request.getClass().getName().toString(), registration.getEndpoint(),
// ((Response) response.getCoapResponse()).getCode(), response.getCode(), request.getPath().toString());
// if (!client.isInit()) {
// client.initReadValue(this.handler, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration));
// }
// /** Not Found */
// if (rpcRequest != null) {
// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), response.getErrorMessage(), LOG_LW2M_ERROR);
// }
// /** Not Found
// set setClient_fw_info... = empty
// **/
// if (client.getFwUpdate() != null && client.getFwUpdate().isInfoFwSwUpdate()) {
// client.getFwUpdate().initReadValue(handler, this, request.getPath().toString());
// }
// if (client.getSwUpdate() != null && client.getSwUpdate().isInfoFwSwUpdate()) {
// client.getSwUpdate().initReadValue(handler, this, request.getPath().toString());
// }
// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) || request.getPath().toString().equals(SW_PACKAGE_ID)) {
// this.afterWriteFwSWUpdateError(registration, request, response.getErrorMessage());
// }
// if (request.getPath().toString().equals(FW_UPDATE_ID) || request.getPath().toString().equals(SW_INSTALL_ID)) {
// this.afterExecuteFwSwUpdateError(registration, request, response.getErrorMessage());
// }
// }
} , e - > {
responseRequestExecutor . submit ( ( ) - > {
callback . onError ( e ) ;
callback . onError ( JacksonUtil . toString ( request ) , e ) ;
} ) ;
// /** version == null
// set setClient_fw_info... = empty
// **/
// if (client.getFwUpdate() != null && client.getFwUpdate().isInfoFwSwUpdate()) {
// client.getFwUpdate().initReadValue(handler, this, request.getPath().toString());
// }
// if (client.getSwUpdate() != null && client.getSwUpdate().isInfoFwSwUpdate()) {
// client.getSwUpdate().initReadValue(handler, this, request.getPath().toString());
// }
// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) || request.getPath().toString().equals(SW_PACKAGE_ID)) {
// this.afterWriteFwSWUpdateError(registration, request, e.getMessage());
// }
// if (request.getPath().toString().equals(FW_UPDATE_ID) || request.getPath().toString().equals(SW_INSTALL_ID)) {
// this.afterExecuteFwSwUpdateError(registration, request, e.getMessage());
// }
// if (!client.isInit()) {
// client.initReadValue(this.handler, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration));
// }
// String msg = String.format("%s: SendRequest %s: Resource path - %s msg error - %s",
// LOG_LW2M_ERROR, request.getClass().getName().toString(), request.getPath().toString(), e.getMessage());
// handler.sendLogsToThingsboard(client, msg);
// log.error("[{}] [{}] - [{}] error SendRequest", request.getClass().getName().toString(), request.getPath().toString(), e.toString());
// if (rpcRequest != null) {
// handler.sentRpcResponse(rpcRequest, CoAP.CodeClass.ERROR_RESPONSE.name(), e.getMessage(), LOG_LW2M_ERROR);
// }
} ) ;
}