@ -18,16 +18,19 @@ package org.thingsboard.server.transport.lwm2m.server.downlink;
import lombok.RequiredArgsConstructor ;
import lombok.RequiredArgsConstructor ;
import lombok.extern.slf4j.Slf4j ;
import lombok.extern.slf4j.Slf4j ;
import org.eclipse.leshan.core.Link ;
import org.eclipse.leshan.core.Link ;
import org.eclipse.leshan.core.LwM2m ;
import org.eclipse.leshan.core.attributes.Attribute ;
import org.eclipse.leshan.core.attributes.Attribute ;
import org.eclipse.leshan.core.attributes.AttributeSet ;
import org.eclipse.leshan.core.attributes.AttributeSet ;
import org.eclipse.leshan.core.model.ObjectModel ;
import org.eclipse.leshan.core.model.ResourceModel ;
import org.eclipse.leshan.core.model.ResourceModel ;
import org.eclipse.leshan.core.node.LwM2mObjectInstance ;
import org.eclipse.leshan.core.node.LwM2mPath ;
import org.eclipse.leshan.core.node.LwM2mPath ;
import org.eclipse.leshan.core.node.LwM2mResource ;
import org.eclipse.leshan.core.node.LwM2mResource ;
import org.eclipse.leshan.core.node.ObjectLink ;
import org.eclipse.leshan.core.node.ObjectLink ;
import org.eclipse.leshan.core.node.codec.CodecException ;
import org.eclipse.leshan.core.observation.Observation ;
import org.eclipse.leshan.core.observation.Observation ;
import org.eclipse.leshan.core.request.CompositeDownlinkRequest ;
import org.eclipse.leshan.core.request.CompositeDownlinkRequest ;
import org.eclipse.leshan.core.request.ContentFormat ;
import org.eclipse.leshan.core.request.ContentFormat ;
import org.eclipse.leshan.core.request.CreateRequest ;
import org.eclipse.leshan.core.request.DeleteRequest ;
import org.eclipse.leshan.core.request.DeleteRequest ;
import org.eclipse.leshan.core.request.DiscoverRequest ;
import org.eclipse.leshan.core.request.DiscoverRequest ;
import org.eclipse.leshan.core.request.DownlinkRequest ;
import org.eclipse.leshan.core.request.DownlinkRequest ;
@ -40,7 +43,9 @@ import org.eclipse.leshan.core.request.WriteAttributesRequest;
import org.eclipse.leshan.core.request.WriteCompositeRequest ;
import org.eclipse.leshan.core.request.WriteCompositeRequest ;
import org.eclipse.leshan.core.request.WriteRequest ;
import org.eclipse.leshan.core.request.WriteRequest ;
import org.eclipse.leshan.core.request.exception.ClientSleepingException ;
import org.eclipse.leshan.core.request.exception.ClientSleepingException ;
import org.eclipse.leshan.core.request.exception.InvalidRequestException ;
import org.eclipse.leshan.core.request.exception.TimeoutException ;
import org.eclipse.leshan.core.request.exception.TimeoutException ;
import org.eclipse.leshan.core.response.CreateResponse ;
import org.eclipse.leshan.core.response.DeleteResponse ;
import org.eclipse.leshan.core.response.DeleteResponse ;
import org.eclipse.leshan.core.response.DiscoverResponse ;
import org.eclipse.leshan.core.response.DiscoverResponse ;
import org.eclipse.leshan.core.response.ExecuteResponse ;
import org.eclipse.leshan.core.response.ExecuteResponse ;
@ -55,7 +60,6 @@ import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.server.model.LwM2mModelProvider ;
import org.eclipse.leshan.server.model.LwM2mModelProvider ;
import org.eclipse.leshan.server.registration.Registration ;
import org.eclipse.leshan.server.registration.Registration ;
import org.springframework.stereotype.Service ;
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.data.device.data.lwm2m.ObjectAttributes ;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent ;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent ;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig ;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig ;
@ -73,9 +77,12 @@ import javax.annotation.PreDestroy;
import java.util.Arrays ;
import java.util.Arrays ;
import java.util.Collection ;
import java.util.Collection ;
import java.util.Date ;
import java.util.Date ;
import java.util.LinkedHashMap ;
import java.util.LinkedList ;
import java.util.LinkedList ;
import java.util.List ;
import java.util.List ;
import java.util.Map ;
import java.util.Set ;
import java.util.Set ;
import java.util.concurrent.ConcurrentHashMap ;
import java.util.function.Function ;
import java.util.function.Function ;
import java.util.function.Predicate ;
import java.util.function.Predicate ;
import java.util.stream.Collectors ;
import java.util.stream.Collectors ;
@ -85,7 +92,11 @@ import static org.eclipse.leshan.core.attributes.Attribute.LESSER_THAN;
import static org.eclipse.leshan.core.attributes.Attribute.MAXIMUM_PERIOD ;
import static org.eclipse.leshan.core.attributes.Attribute.MAXIMUM_PERIOD ;
import static org.eclipse.leshan.core.attributes.Attribute.MINIMUM_PERIOD ;
import static org.eclipse.leshan.core.attributes.Attribute.MINIMUM_PERIOD ;
import static org.eclipse.leshan.core.attributes.Attribute.STEP ;
import static org.eclipse.leshan.core.attributes.Attribute.STEP ;
import static org.eclipse.leshan.core.model.ResourceModel.Type.OBJLNK ;
import static org.eclipse.leshan.core.model.ResourceModel.Type.OPAQUE ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertMultiResourceValuesFromRpcBody ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validateVersionedId ;
@Slf4j
@Slf4j
@Service
@Service
@ -123,40 +134,49 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
@Override
@Override
public void sendReadRequest ( LwM2mClient client , TbLwM2MReadRequest request , DownlinkRequestCallback < ReadRequest , ReadResponse > callback ) {
public void sendReadRequest ( LwM2mClient client , TbLwM2MReadRequest request , DownlinkRequestCallback < ReadRequest , ReadResponse > callback ) {
validateVersionedId ( client , request ) ;
try {
ReadRequest downlink = new ReadRequest ( getRequestContentFormat ( client , request , this . config . getModelProvider ( ) ) , request . getObjectId ( ) ) ;
validateVersionedId ( client , request ) ;
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
ReadRequest downlink = new ReadRequest ( getReadRequestContentFormat ( client , request , this . config . getModelProvider ( ) ) , request . getObjectId ( ) ) ;
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} catch ( InvalidRequestException e ) {
callback . onValidationError ( request . toString ( ) , e . getMessage ( ) ) ;
}
}
}
@Override
@Override
public void sendReadCompositeRequest ( LwM2mClient client , TbLwM2MReadCompositeRequest request , DownlinkRequestCallback < ReadCompositeRequest , ReadCompositeResponse > callback ) {
public void sendReadCompositeRequest ( LwM2mClient client , TbLwM2MReadCompositeRequest request ,
validateVersionedIds ( client , request ) ;
DownlinkRequestCallback < ReadCompositeRequest , ReadCompositeResponse > callback , ContentFormat contentFormatComposite ) {
ContentFormat requestContentFormat = ContentFormat . SENML_JSON ;
try {
ContentFormat responseContentFormat = ContentFormat . SENML_JSON ;
ReadCompositeRequest downlink = new ReadCompositeRequest ( contentFormatComposite , contentFormatComposite , request . getObjectIds ( ) ) ;
sendCompositeRequest ( client , downlink , this . config . getTimeout ( ) , callback ) ;
ReadCompositeRequest downlink = new ReadCompositeRequest ( requestContentFormat , responseContentFormat , request . getObjectIds ( ) ) ;
} catch ( InvalidRequestException e ) {
sendCompositeRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
callback . onValidationError ( request . toString ( ) , e . getMessage ( ) ) ;
}
}
}
@Override
@Override
public void sendObserveRequest ( LwM2mClient client , TbLwM2MObserveRequest request , DownlinkRequestCallback < ObserveRequest , ObserveResponse > callback ) {
public void sendObserveRequest ( LwM2mClient client , TbLwM2MObserveRequest request , DownlinkRequestCallback < ObserveRequest , ObserveResponse > callback ) {
validateVersionedId ( client , request ) ;
try {
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
validateVersionedId ( client , request ) ;
Set < Observation > observations = context . getServer ( ) . getObservationService ( ) . getObservations ( client . getRegistration ( ) ) ;
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
if ( observations . stream ( ) . noneMatch ( observation - > observation . getPath ( ) . equals ( resultIds ) ) ) {
Set < Observation > observations = context . getServer ( ) . getObservationService ( ) . getObservations ( client . getRegistration ( ) ) ;
ObserveRequest downlink ;
if ( observations . stream ( ) . noneMatch ( observation - > observation . getPath ( ) . equals ( resultIds ) ) ) {
ContentFormat contentFormat = getRequestContentFormat ( client , request , this . config . getModelProvider ( ) ) ;
ObserveRequest downlink ;
if ( resultIds . isResource ( ) ) {
ContentFormat contentFormat = getReadRequestContentFormat ( client , request , this . config . getModelProvider ( ) ) ;
downlink = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) ) ;
if ( resultIds . isResource ( ) ) {
} else if ( resultIds . isObjectInstance ( ) ) {
downlink = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) ) ;
downlink = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) ) ;
} else if ( resultIds . isObjectInstance ( ) ) {
downlink = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) ) ;
} else {
downlink = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) ) ;
}
log . info ( "[{}] Send observation: {}." , client . getEndpoint ( ) , request . getVersionedId ( ) ) ;
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} else {
} else {
downlink = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) ) ;
callback . onValidationError ( resultIds . toString ( ) , "Observation is already registered!" ) ;
}
}
log . info ( "[{}] Send observation: {}." , client . getEndpoint ( ) , request . getVersionedId ( ) ) ;
} catch ( InvalidRequestException e ) {
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
callback . onValidationError ( request . toString ( ) , e . getMessage ( ) ) ;
} else {
callback . onValidationError ( resultIds . toString ( ) , "Observation is already registered!" ) ;
}
}
}
}
@ -174,25 +194,36 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
@Override
@Override
public void sendExecuteRequest ( LwM2mClient client , TbLwM2MExecuteRequest request , DownlinkRequestCallback < ExecuteRequest , ExecuteResponse > callback ) {
public void sendExecuteRequest ( LwM2mClient client , TbLwM2MExecuteRequest request , DownlinkRequestCallback < ExecuteRequest , ExecuteResponse > callback ) {
ResourceModel resourceModelExecute = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
try {
if ( resourceModelExecute ! = null ) {
ResourceModel resourceModelExecute = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
ExecuteRequest downlink ;
if ( resourceModelExecute ! = null ) {
if ( request . getParams ( ) ! = null & & ! resourceModelExecute . multiple ) {
validateVersionedId ( client , request ) ;
downlink = new ExecuteRequest ( request . getVersionedId ( ) , ( String ) this . converter . convertValue ( request . getParams ( ) , resourceModelExecute . type , ResourceModel . Type . STRING , new LwM2mPath ( request . getObjectId ( ) ) ) ) ;
ExecuteRequest downlink ;
} else {
if ( request . getParams ( ) ! = null & & ! resourceModelExecute . multiple ) {
downlink = new ExecuteRequest ( request . getVersionedId ( ) ) ;
downlink = new ExecuteRequest ( request . getObjectId ( ) , ( String ) this . converter . convertValue ( request . getParams ( ) , resourceModelExecute . type , ResourceModel . Type . STRING , new LwM2mPath ( request . getObjectId ( ) ) ) ) ;
} else {
downlink = new ExecuteRequest ( request . getObjectId ( ) ) ;
}
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
}
}
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} catch ( InvalidRequestException e ) {
callback . onValidationError ( request . toString ( ) , e . getMessage ( ) ) ;
}
}
}
}
@Override
@Override
public void sendDeleteRequest ( LwM2mClient client , TbLwM2MDeleteRequest request , DownlinkRequestCallback < DeleteRequest , DeleteResponse > callback ) {
public void sendDeleteRequest ( LwM2mClient client , TbLwM2MDeleteRequest request , DownlinkRequestCallback < DeleteRequest , DeleteResponse > callback ) {
sendSimpleRequest ( client , new DeleteRequest ( request . getObjectId ( ) ) , request . getTimeout ( ) , callback ) ;
try {
validateVersionedId ( client , request ) ;
sendSimpleRequest ( client , new DeleteRequest ( request . getObjectId ( ) ) , request . getTimeout ( ) , callback ) ;
} catch ( InvalidRequestException e ) {
callback . onValidationError ( request . toString ( ) , e . getMessage ( ) ) ;
}
}
}
@Override
@Override
public void sendCancelObserveRequest ( LwM2mClient client , TbLwM2MCancelObserveRequest request , DownlinkRequestCallback < TbLwM2MCancelObserveRequest , Integer > callback ) {
public void sendCancelObserveRequest ( LwM2mClient client , TbLwM2MCancelObserveRequest request , DownlinkRequestCallback < TbLwM2MCancelObserveRequest , Integer > callback ) {
validateVersionedId ( client , request ) ;
int observeCancelCnt = context . getServer ( ) . getObservationService ( ) . cancelObservations ( client . getRegistration ( ) , request . getObjectId ( ) ) ;
int observeCancelCnt = context . getServer ( ) . getObservationService ( ) . cancelObservations ( client . getRegistration ( ) , request . getObjectId ( ) ) ;
callback . onSuccess ( request , observeCancelCnt ) ;
callback . onSuccess ( request , observeCancelCnt ) ;
}
}
@ -209,51 +240,86 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
sendSimpleRequest ( client , new DiscoverRequest ( request . getObjectId ( ) ) , request . getTimeout ( ) , callback ) ;
sendSimpleRequest ( client , new DiscoverRequest ( request . getObjectId ( ) ) , request . getTimeout ( ) , callback ) ;
}
}
/ * *
* Example # 1 :
* AttributeSet attributes = new AttributeSet ( new Attribute ( Attribute . MINIMUM_PERIOD , 10L ) ,
* new Attribute ( Attribute . MAXIMUM_PERIOD , 100L ) ) ;
* WriteAttributesRequest requestTest = new WriteAttributesRequest ( 3 , 0 , 14 , attributes ) ;
* sendSimpleRequest ( client , requestTest , request . getTimeout ( ) , callback ) ;
* < p >
* Example # 2
* Dimension and Object version are read only attributes .
* addAttribute ( attributes , DIMENSION , params . getDim ( ) , dim - > dim > = 0 & & dim < = 255 ) ;
* addAttribute ( attributes , OBJECT_VERSION , params . getVer ( ) , StringUtils : : isNotEmpty , Function . identity ( ) ) ;
* /
@Override
@Override
public void sendWriteAttributesRequest ( LwM2mClient client , TbLwM2MWriteAttributesRequest request , DownlinkRequestCallback < WriteAttributesRequest , WriteAttributesResponse > callback ) {
public void sendWriteAttributesRequest ( LwM2mClient client , TbLwM2MWriteAttributesRequest request , DownlinkRequestCallback < WriteAttributesRequest , WriteAttributesResponse > callback ) {
validateVersionedId ( client , request ) ;
try {
if ( request . getAttributes ( ) = = null ) {
validateVersionedId ( client , request ) ;
throw new IllegalArgumentException ( "Attributes to write are not specified!" ) ;
if ( request . getAttributes ( ) = = null ) {
throw new IllegalArgumentException ( "Attributes to write are not specified!" ) ;
}
ObjectAttributes params = request . getAttributes ( ) ;
List < Attribute > attributes = new LinkedList < > ( ) ;
addAttribute ( attributes , MAXIMUM_PERIOD , params . getPmax ( ) ) ;
addAttribute ( attributes , MINIMUM_PERIOD , params . getPmin ( ) ) ;
addAttribute ( attributes , GREATER_THAN , params . getGt ( ) ) ;
addAttribute ( attributes , LESSER_THAN , params . getLt ( ) ) ;
addAttribute ( attributes , STEP , params . getSt ( ) ) ;
AttributeSet attributeSet = new AttributeSet ( attributes ) ;
sendSimpleRequest ( client , new WriteAttributesRequest ( request . getObjectId ( ) , attributeSet ) , request . getTimeout ( ) , callback ) ;
} catch ( InvalidRequestException e ) {
callback . onValidationError ( request . toString ( ) , e . getMessage ( ) ) ;
}
}
ObjectAttributes params = request . getAttributes ( ) ;
List < Attribute > attributes = new LinkedList < > ( ) ;
// Dimension and Object version are read only attributes.
// addAttribute(attributes, DIMENSION, params.getDim(), dim -> dim >= 0 && dim <= 255);
// addAttribute(attributes, OBJECT_VERSION, params.getVer(), StringUtils::isNotEmpty, Function.identity());
addAttribute ( attributes , MAXIMUM_PERIOD , params . getPmax ( ) ) ;
addAttribute ( attributes , MINIMUM_PERIOD , params . getPmin ( ) ) ;
addAttribute ( attributes , GREATER_THAN , params . getGt ( ) ) ;
addAttribute ( attributes , LESSER_THAN , params . getLt ( ) ) ;
addAttribute ( attributes , STEP , params . getSt ( ) ) ;
AttributeSet attributeSet = new AttributeSet ( attributes ) ;
sendSimpleRequest ( client , new WriteAttributesRequest ( request . getObjectId ( ) , attributeSet ) , request . getTimeout ( ) , callback ) ;
}
}
@Override
@Override
public void sendWriteReplaceRequest ( LwM2mClient client , TbLwM2MWriteReplaceRequest request , DownlinkRequestCallback < WriteRequest , WriteResponse > callback ) {
public void sendWriteReplaceRequest ( LwM2mClient client , TbLwM2MWriteReplaceRequest request , DownlinkRequestCallback < WriteRequest , WriteResponse > callback ) {
ResourceModel resourceModelWrite = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
if ( resourceModelWrite ! = null ) {
if ( resultIds . isResource ( ) | | resultIds . isResourceInstance ( ) ) {
ContentFormat contentFormat = convertResourceModelTypeToContentFormat ( client , resourceModelWrite . type ) ;
validateVersionedId ( client , request ) ;
try {
ResourceModel resourceModelWrite = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
LwM2mPath path = new LwM2mPath ( request . getObjectId ( ) ) ;
if ( resourceModelWrite ! = null ) {
WriteRequest downlink = this . getWriteRequestSingleResource ( resourceModelWrite . type , contentFormat ,
ContentFormat contentFormat = getWriteRequestContentFormat ( client , request , this . config . getModelProvider ( ) ) ;
path . getObjectId ( ) , path . getObjectInstanceId ( ) , path . getResourceId ( ) , request . getValue ( ) ) ;
try {
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
WriteRequest downlink = null ;
} catch ( Exception e ) {
if ( resourceModelWrite . multiple ) {
callback . onError ( toString ( request ) , e ) ;
if ( request . getValue ( ) instanceof Map & & ( ( Map ) request . getValue ( ) ) . size ( ) > 0 ) {
downlink = new WriteRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) ,
( Map < Integer , ? > ) request . getValue ( ) , resourceModelWrite . type ) ;
} else {
callback . onValidationError ( toString ( request ) , "Resource value is: " + request . getValue ( ) . getClass ( ) . getSimpleName ( ) + ". Value of Multi-Instance Resource must be in Json format!" ) ;
}
} else {
downlink = this . getWriteRequestSingleResource ( resourceModelWrite . type , contentFormat ,
resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) , request . getValue ( ) ) ;
}
if ( downlink ! = null ) {
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} else {
callback . onValidationError ( toString ( request ) , "WriteRequest is null." ) ;
}
} catch ( Exception e ) {
callback . onError ( toString ( request ) , e ) ;
}
} else {
callback . onValidationError ( toString ( request ) , "Resource " + request . getVersionedId ( ) + " is not configured in the device profile!" ) ;
}
}
} else {
} else {
callback . onValidationError ( toString ( request ) , "Resource " + request . getVersionedId ( ) + " is not configured in the device profile!" ) ;
callback . onValidationError ( toString ( request ) , "Resource " + request . getVersionedId ( ) + ". This operation can only be used for Resource or ResourceInstanc e!" ) ;
}
}
}
}
@Override
@Override
public void sendWriteCompositeRequest ( LwM2mClient client , RpcWriteCompositeRequest rpcWriteCompositeRequest , DownlinkRequestCallback < WriteCompositeRequest , WriteCompositeResponse > callback ) {
public void sendWriteCompositeRequest ( LwM2mClient client , RpcWriteCompositeRequest rpcWriteCompositeRequest ,
ContentFormat contentFormat = ContentFormat . SENML_JSON ;
DownlinkRequestCallback < WriteCompositeRequest , WriteCompositeResponse > callback , ContentFormat contentFormatComposite ) {
try {
try {
WriteCompositeRequest downlink = new WriteCompositeRequest ( contentFormat , rpcWriteCompositeRequest . getNodes ( ) ) ;
WriteCompositeRequest downlink = new WriteCompositeRequest ( contentFormatComposite , rpcWriteCompositeRequest . getNodes ( ) ) ;
//TODO: replace config.getTimeout();
//TODO: replace config.getTimeout();
sendWriteCompositeRequest ( client , downlink , config . getTimeout ( ) , callback ) ;
sendWriteCompositeRequest ( client , downlink , this . config . getTimeout ( ) , callback ) ;
} catch ( InvalidRequestException e ) {
callback . onValidationError ( rpcWriteCompositeRequest . toString ( ) , e . getMessage ( ) ) ;
} catch ( Exception e ) {
} catch ( Exception e ) {
callback . onError ( toString ( rpcWriteCompositeRequest ) , e ) ;
callback . onError ( toString ( rpcWriteCompositeRequest ) , e ) ;
}
}
@ -261,38 +327,90 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
@Override
@Override
public void sendWriteUpdateRequest ( LwM2mClient client , TbLwM2MWriteUpdateRequest request , DownlinkRequestCallback < WriteRequest , WriteResponse > callback ) {
public void sendWriteUpdateRequest ( LwM2mClient client , TbLwM2MWriteUpdateRequest request , DownlinkRequestCallback < WriteRequest , WriteResponse > callback ) {
try {
validateVersionedId ( client , request ) ;
WriteRequest downlink = null ;
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
ContentFormat contentFormat = getWriteRequestContentFormat ( client , request , this . config . getModelProvider ( ) ) ;
if ( resultIds . isObjectInstance ( ) ) {
/ *
* params = "{\"id\":0,\"value\":[{\"id\":14,\"value\":\"+5\"},{\"id\":15,\"value\":\"+9\"}]}"
* int rscId = resultIds . getObjectInstanceId ( ) ;
* contentFormat – Format of the payload ( TLV or JSON ) .
* /
Collection < LwM2mResource > resources = client . getNewResourcesForInstance ( request . getVersionedId ( ) , request . getValue ( ) , this . config . getModelProvider ( ) , this . converter ) ;
if ( resources . size ( ) > 0 ) {
downlink = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resources ) ;
} else {
callback . onValidationError ( toString ( request ) , "No resources to update!" ) ;
}
} else if ( resultIds . isResource ( ) ) {
ResourceModel resourceModelWrite = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
if ( resourceModelWrite . multiple ) {
if ( request . getValue ( ) instanceof Map & & ( ( Map ) request . getValue ( ) ) . size ( ) > 0 ) {
Map value = convertMultiResourceValuesFromRpcBody ( ( LinkedHashMap ) request . getValue ( ) , resourceModelWrite . type , request . getObjectId ( ) ) ;
downlink = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) ,
value , resourceModelWrite . type ) ;
} else {
callback . onValidationError ( toString ( request ) , "Resource value is bad. Format: " + request . getValue ( ) . getClass ( ) . getSimpleName ( ) + ". Value of Multi-Instance Resource must be in Json format!" ) ;
}
}
}
if ( downlink ! = null ) {
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} else {
callback . onValidationError ( toString ( request ) , "Resource " + request . getVersionedId ( ) + ". This operation can only be used for ObjectInstance or Multi-Instance Resource !" ) ;
}
} catch ( Exception e ) {
callback . onValidationError ( toString ( request ) , e . getMessage ( ) ) ;
}
}
public void sendCreateRequest ( LwM2mClient client , TbLwM2MCreateRequest request , DownlinkRequestCallback < CreateRequest , CreateResponse > callback ) {
validateVersionedId ( client , request ) ;
CreateRequest downlink = null ;
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
LwM2mPath resultIds = new LwM2mPath ( request . getObjectId ( ) ) ;
if ( resultIds . isResource ( ) ) {
ObjectModel objectModel = client . getObjectModel ( request . getObjectId ( ) , this . config . getModelProvider ( ) ) ;
/ *
// POST /{Object ID}/{Object Instance ID} && Resources is Mandatory
* send request : path = ' / 3 / 0 ' node = = wM2mObjectInstance
if ( objectModel . multiple ) {
* with params = = "\"resources\" : { 15 : resource : { id : 15 . value : ' + 01 ' . . . } }
// LwM2M CBOR, SenML CBOR, SenML JSON, or TLV (see [LwM2M-CORE])
* * /
ContentFormat contentFormat = getWriteRequestContentFormat ( client , request , this . config . getModelProvider ( ) ) ;
Collection < LwM2mResource > resources = client . getNewResourceForInstance ( request . getVersionedId ( ) , request . getValue ( ) , this . config . getModelProvider ( ) , this . converter ) ;
if ( resultIds . isObject ( ) | | resultIds . isObjectInstance ( ) ) {
ResourceModel resourceModelWrite = client . getResourceModel ( request . getVersionedId ( ) , this . config . getModelProvider ( ) ) ;
Collection < LwM2mResource > resources ;
ContentFormat contentFormat = request . getObjectContentFormat ( ) ! = null ? request . getObjectContentFormat ( ) : convertResourceModelTypeToContentFormat ( client , resourceModelWrite . type ) ;
if ( resultIds . isObject ( ) ) {
WriteRequest downlink = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) ,
// contentFormat = ContentFormat.TLV;
resultIds . getObjectInstanceId ( ) , resources ) ;
if ( request . getValue ( ) ! = null ) {
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
resources = client . getNewResourcesForInstance ( request . getVersionedId ( ) , request . getValue ( ) , this . config . getModelProvider ( ) , this . converter ) ;
} else if ( resultIds . isObjectInstance ( ) ) {
downlink = new CreateRequest ( contentFormat , resultIds . getObjectId ( ) , resources ) ;
/ *
} else if ( request . getNodes ( ) ! = null & & request . getNodes ( ) . size ( ) > 0 ) {
* params = "{\"id\":0,\"resources\":[{\"id\":14,\"value\":\"+5\"},{\"id\":15,\"value\":\"+9\"}]}"
Set < LwM2mObjectInstance > instances = ConcurrentHashMap . newKeySet ( ) ;
* int rscId = resultIds . getObjectInstanceId ( ) ;
request . getNodes ( ) . forEach ( ( key , value ) - > {
* contentFormat – Format of the payload ( TLV or JSON ) .
Collection < LwM2mResource > resourcesForInstance = client . getNewResourcesForInstance ( request . getVersionedId ( ) , value , this . config . getModelProvider ( ) , this . converter ) ;
* /
LwM2mObjectInstance instance = new LwM2mObjectInstance ( Integer . parseInt ( key ) , resourcesForInstance ) ;
Collection < LwM2mResource > resources = client . getNewResourcesForInstance ( request . getVersionedId ( ) , request . getValue ( ) , this . config . getModelProvider ( ) , this . converter ) ;
instances . add ( instance ) ;
if ( resources . size ( ) > 0 ) {
} ) ;
ContentFormat contentFormat = request . getObjectContentFormat ( ) ! = null ? request . getObjectContentFormat ( ) : ContentFormat . DEFAULT ;
LwM2mObjectInstance [ ] instanceArrays = instances . toArray ( new LwM2mObjectInstance [ instances . size ( ) ] ) ;
WriteRequest downlink = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resources ) ;
downlink = new CreateRequest ( contentFormat , resultIds . getObjectId ( ) , instanceArrays ) ;
}
} else {
resources = client . getNewResourcesForInstance ( request . getVersionedId ( ) , request . getValue ( ) , this . config . getModelProvider ( ) , this . converter ) ;
LwM2mObjectInstance instance = new LwM2mObjectInstance ( resultIds . getObjectInstanceId ( ) , resources ) ;
downlink = new CreateRequest ( contentFormat , resultIds . getObjectId ( ) , instance ) ;
}
}
if ( downlink ! = null ) {
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
sendSimpleRequest ( client , downlink , request . getTimeout ( ) , callback ) ;
} else {
} else {
callback . onValidationError ( toString ( request ) , "No resources to update!" ) ;
callback . onValidationError ( toString ( request ) , "Path " + request . getVersionedId ( ) +
". This operation can only be used for created new ObjectInstance !" ) ;
}
}
} else {
} else {
callback . onValidationError ( toString ( request ) , "Update of the root level object is not supported yet!" ) ;
throw new IllegalArgumentException ( "Path " + request . getVersionedId ( ) +
". Object must be Multiple !" ) ;
}
}
}
}
private < R extends SimpleDownlinkRequest < T > , T extends LwM2mResponse > void sendSimpleRequest ( LwM2mClient client , R request , long timeoutInMs , DownlinkRequestCallback < R , T > callback ) {
private < R extends SimpleDownlinkRequest < T > , T extends LwM2mResponse > void sendSimpleRequest ( LwM2mClient client , R request , long timeoutInMs , DownlinkRequestCallback < R , T > callback ) {
sendRequest ( client , request , timeoutInMs , callback , r - > request . getPath ( ) . toString ( ) ) ;
sendRequest ( client , request , timeoutInMs , callback , r - > request . getPath ( ) . toString ( ) ) ;
}
}
@ -301,6 +419,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
sendRequest ( client , request , timeoutInMs , callback , r - > request . getPaths ( ) . toString ( ) ) ;
sendRequest ( client , request , timeoutInMs , callback , r - > request . getPaths ( ) . toString ( ) ) ;
}
}
private < R extends DownlinkRequest < T > , T extends LwM2mResponse > void sendRequest ( LwM2mClient client , R request , long timeoutInMs , DownlinkRequestCallback < R , T > callback , Function < R , String > pathToStringFunction ) {
private < R extends DownlinkRequest < T > , T extends LwM2mResponse > void sendRequest ( LwM2mClient client , R request , long timeoutInMs , DownlinkRequestCallback < R , T > callback , Function < R , String > pathToStringFunction ) {
if ( ! clientContext . isDownlinkAllowed ( client ) ) {
if ( ! clientContext . isDownlinkAllowed ( client ) ) {
log . trace ( "[{}] ignore downlink request cause client is sleeping." , client . getEndpoint ( ) ) ;
log . trace ( "[{}] ignore downlink request cause client is sleeping." , client . getEndpoint ( ) ) ;
@ -319,7 +438,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
clientContext . awake ( client ) ;
clientContext . awake ( client ) ;
}
}
} ) ;
} ) ;
} , e - > handleDownlinkError ( client , request , callback , e ) ) ;
} , e - > handleDownlinkError ( client , request , callback , e ) ) ;
} catch ( Exception e ) {
} catch ( Exception e ) {
handleDownlinkError ( client , request , callback , e ) ;
handleDownlinkError ( client , request , callback , e ) ;
}
}
@ -366,6 +485,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
} ) ;
} ) ;
}
}
private WriteRequest getWriteRequestSingleResource ( ResourceModel . Type type , ContentFormat contentFormat , int objectId , int instanceId , int resourceId , Object value ) {
private WriteRequest getWriteRequestSingleResource ( ResourceModel . Type type , ContentFormat contentFormat , int objectId , int instanceId , int resourceId , Object value ) {
switch ( type ) {
switch ( type ) {
case STRING : // String
case STRING : // String
@ -395,24 +515,6 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
}
}
}
}
private void validateVersionedId ( LwM2mClient client , HasVersionedId request ) {
client . isValidObjectVersion ( request . getVersionedId ( ) ) ;
if ( request . getObjectId ( ) = = null ) {
throw new IllegalArgumentException ( "Specified object id is null!" ) ;
}
}
private void validateVersionedIds ( LwM2mClient client , HasVersionedIds request ) {
for ( String versionedId : request . getVersionedIds ( ) ) {
client . isValidObjectVersion ( versionedId ) ;
}
for ( String objectId : request . getObjectIds ( ) ) {
if ( objectId = = null ) {
throw new IllegalArgumentException ( "Specified object id is null!" ) ;
}
}
}
private static < T > void addAttribute ( List < Attribute > attributes , String attributeName , T value ) {
private static < T > void addAttribute ( List < Attribute > attributes , String attributeName , T value ) {
addAttribute ( attributes , attributeName , value , null , null ) ;
addAttribute ( attributes , attributeName , value , null , null ) ;
}
}
@ -427,52 +529,91 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im
}
}
}
}
private static ContentFormat convertResourceModelTypeToContentFormat ( LwM2mClient client , ResourceModel . Type type ) {
private static < T extends HasContentFormat & HasVersionedId > ContentFormat getReadRequestContentFormat ( LwM2mClient client , T request , LwM2mModelProvider modelProvider ) {
switch ( type ) {
if ( request . getRequestContentFormat ( ) . isPresent ( ) ) {
case BOOLEAN :
return request . getRequestContentFormat ( ) . get ( ) ;
case STRING :
} else {
case TIME :
return getRequestContentFormat ( client , request . getVersionedId ( ) , modelProvider ) ;
case INTEGER :
case FLOAT :
return client . getDefaultContentFormat ( ) ;
case OPAQUE :
return ContentFormat . OPAQUE ;
case OBJLNK :
return ContentFormat . LINK ;
default :
}
}
throw new CodecException ( "Invalid ResourceModel_Type for %s ContentFormat." , type ) ;
}
}
private static ContentFormat getRequestContentFormat ( LwM2mClient client , HasContentFormat request , LwM2mModelProvider modelProvider ) {
private static ContentFormat getWriteRequestContentFormat ( LwM2mClient client , TbLwM2MDownlinkRequest request , LwM2mModelProvider modelProvider ) {
if ( request . getRequestContentFormat ( ) ! = null ) {
if ( request instanceof TbLwM2MWriteReplaceRequest & & ( ( TbLwM2MWriteReplaceRequest ) request ) . getContentFormat ( ) ! = null ) {
return request . getRequestContentFormat ( ) ;
return ( ( TbLwM2MWriteReplaceRequest ) request ) . getContentFormat ( ) ;
} else if ( request instanceof TbLwM2MWriteUpdateRequest & & ( ( TbLwM2MWriteUpdateRequest ) request ) . getObjectContentFormat ( ) ! = null ) {
return ( ( TbLwM2MWriteUpdateRequest ) request ) . getObjectContentFormat ( ) ;
} else {
} else {
String versionedId = null ;
String versionedId = null ;
if ( request instanceof TbLwM2MReadRequest ) {
if ( request instanceof TbLwM2MWriteReplaceRequest ) {
versionedId = ( ( TbLwM2MReadRequest ) request ) . getVersionedId ( ) ;
versionedId = ( ( TbLwM2MWriteReplaceRequest ) request ) . getVersionedId ( ) ;
} else if ( request instanceof TbLwM2MObserveRequest ) {
} else if ( request instanceof TbLwM2MWriteUpdateRequest ) {
versionedId = ( ( TbLwM2MObserveRequest ) request ) . getVersionedId ( ) ;
versionedId = ( ( TbLwM2MWriteUpdateRequest ) request ) . getVersionedId ( ) ;
} else if ( request instanceof TbLwM2MCreateRequest ) {
versionedId = ( ( TbLwM2MCreateRequest ) request ) . getVersionedId ( ) ;
}
return getRequestContentFormat ( client , versionedId , modelProvider ) ;
}
}
private static ContentFormat getRequestContentFormat ( LwM2mClient client , String versionedId , LwM2mModelProvider modelProvider ) {
LwM2mPath pathIds = new LwM2mPath ( fromVersionedIdToObjectId ( versionedId ) ) ;
if ( pathIds . isResourceInstance ( ) ) {
ResourceModel resourceModel = client . getResourceModel ( versionedId , modelProvider ) ;
if ( OBJLNK . equals ( resourceModel . type ) ) {
return ContentFormat . LINK ;
} else if ( OPAQUE . equals ( resourceModel . type ) ) {
return ContentFormat . OPAQUE ;
} else {
return findFirst ( client . getClientSupportContentFormats ( ) , client . getDefaultContentFormat ( ) , ContentFormat . CBOR , ContentFormat . SENML_CBOR , ContentFormat . SENML_JSON ) ;
}
} else if ( pathIds . isResource ( ) ) {
ResourceModel resourceModel = client . getResourceModel ( versionedId , modelProvider ) ;
if ( ! resourceModel . multiple ) {
if ( OBJLNK . equals ( resourceModel . type ) ) {
return ContentFormat . LINK ;
} else if ( OPAQUE . equals ( resourceModel . type ) ) {
return ContentFormat . OPAQUE ;
} else {
return findFirst ( client . getClientSupportContentFormats ( ) , client . getDefaultContentFormat ( ) , ContentFormat . CBOR , ContentFormat . SENML_CBOR , ContentFormat . SENML_JSON ) ;
}
} else {
return getContentFormatForComplex ( client ) ;
}
}
String id = fromVersionedIdToObjectId ( versionedId ) ;
} else {
if ( id ! = null & & new LwM2mPath ( id ) . isResource ( ) & & ! client . isResourceMultiInstances ( versionedId , modelProvider ) ) {
return getContentFormatForComplex ( client ) ;
return client . getDefaultContentFormat ( ) ;
}
}
private static ContentFormat getContentFormatForComplex ( LwM2mClient client ) {
if ( LwM2m . Version . V1_0 . equals ( client . getRegistration ( ) . getLwM2mVersion ( ) ) ) {
return ContentFormat . TLV ;
} else if ( LwM2m . Version . V1_1 . equals ( client . getRegistration ( ) . getLwM2mVersion ( ) ) ) {
ContentFormat result = findFirst ( client . getClientSupportContentFormats ( ) , null , ContentFormat . SENML_CBOR , ContentFormat . SENML_JSON , ContentFormat . TLV , ContentFormat . JSON ) ;
if ( result ! = null ) {
return result ;
} else {
} else {
return ContentFormat . DEFAULT ;
throw new RuntimeException ( "The client does not support any of SenML CBOR, SenML JSON, TLV or JSON formats. Can't send complex requests. Try using singe-instance requests." ) ;
}
} else {
throw new RuntimeException ( "The version " + client . getRegistration ( ) . getLwM2mVersion ( ) + " is not supported!" ) ;
}
}
private static ContentFormat findFirst ( Set < ContentFormat > supported , ContentFormat defaultValue , ContentFormat . . . desiredFormats ) {
for ( ContentFormat contentFormat : desiredFormats ) {
if ( supported . contains ( contentFormat ) ) {
return contentFormat ;
}
}
}
}
return defaultValue ;
}
}
private < R > String toString ( R request ) {
private < R > String toString ( R request ) {
try {
try {
try {
return request ! = null ? request . toString ( ) : "" ;
return JacksonUtil . toString ( request ) ;
} catch ( Exception e ) {
return request . toString ( ) ;
}
} catch ( Exception e ) {
} catch ( Exception e ) {
log . warn ( "Failed to convert request to string" , e ) ;
log . trace ( "Failed to convert request to string" , e ) ;
return request ! = null ? request . getClass ( ) . getSimpleName ( ) : "" ;
return request . getClass ( ) . getSimpleName ( ) ;
}
}
}
}
}
}