@ -19,6 +19,7 @@ import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j ;
import org.eclipse.californium.core.coap.CoAP ;
import org.eclipse.californium.core.coap.Response ;
import org.eclipse.leshan.core.Link ;
import org.eclipse.leshan.core.model.ResourceModel ;
import org.eclipse.leshan.core.node.LwM2mNode ;
import org.eclipse.leshan.core.node.LwM2mPath ;
@ -35,7 +36,6 @@ import org.eclipse.leshan.core.request.ObserveRequest;
import org.eclipse.leshan.core.request.ReadRequest ;
import org.eclipse.leshan.core.request.WriteRequest ;
import org.eclipse.leshan.core.request.exception.ClientSleepingException ;
import org.eclipse.leshan.core.response.CancelObservationResponse ;
import org.eclipse.leshan.core.response.DeleteResponse ;
import org.eclipse.leshan.core.response.DiscoverResponse ;
import org.eclipse.leshan.core.response.ExecuteResponse ;
@ -68,15 +68,26 @@ import java.util.stream.Collectors;
import static org.eclipse.californium.core.coap.CoAP.ResponseCode.CONTENT ;
import static org.eclipse.leshan.core.ResponseCode.BAD_REQUEST ;
import static org.eclipse.leshan.core.ResponseCode.NOT_FOUND ;
import static org.thingsboard.server.common.data.firmware.FirmwareUpdateStatus.DOWNLOADED ;
import static org.thingsboard.server.common.data.firmware.FirmwareUpdateStatus.FAILED ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getContentFormatByResourceModelType ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.DEFAULT_TIMEOUT ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FR_PATH_RESOURCE_VER_ID ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_PACKAGE_ID ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_UPDATE_ID ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_ERROR ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_INFO ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_VALUE ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.DISCOVER ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.DISCOVER_All ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.EXECUTE ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_CANCEL ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.OBSERVE_READ_ALL ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.RESPONSE_CHANNEL ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_REPLACE ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_UPDATE ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.RESPONSE_REQUEST_CHANNEL ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SW_INSTALL_ID ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SW_PACKAGE_ID ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertPathFromIdVerToObjectId ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertPathFromObjectIdToIdVer ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.createWriteAttributeRequest ;
@ -86,20 +97,20 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.c
@TbLwM2mTransportComponent
@RequiredArgsConstructor
public class LwM2mTransportRequest {
private ExecutorService executorResponse ;
private ExecutorService responseRequ estE xecutor;
public LwM2mValueConverterImpl converter ;
private final LwM2mTransportContext context ;
private final LwM2MTransportServerConfig config ;
private final LwM2mClientContext lwM2mClientContext ;
private final DefaultLwM2MTransportMsgHandler serviceImpl ;
private final DefaultLwM2MTransportMsgHandler handler ;
@PostConstruct
public void init ( ) {
this . converter = LwM2mValueConverterImpl . getInstance ( ) ;
executorResponse = Executors . newFixedThreadPool ( this . config . getResponsePoolSize ( ) ,
new NamedThreadFactory ( String . format ( "LwM2M %s channel response" , RESPONSE_CHANNEL ) ) ) ;
responseRequ estE xecutor = Executors . newFixedThreadPool ( this . config . getResponsePoolSize ( ) ,
new NamedThreadFactory ( String . format ( "LwM2M %s channel response after request " , RESPONSE_REQUEST _CHANNEL ) ) ) ;
}
/ * *
@ -115,154 +126,186 @@ public class LwM2mTransportRequest {
String contentFormatName , Object params , long timeoutInMs , Lwm2mClientRpcRequest rpcRequest ) {
try {
String target = convertPathFromIdVerToObjectId ( targetIdVer ) ;
DownlinkRequest request = null ;
ContentFormat contentFormat = contentFormatName ! = null ? ContentFormat . fromName ( contentFormatName . toUpperCase ( ) ) : ContentFormat . DEFAULT ;
LwM2mClient lwM2MClient = this . lwM2mClientContext . getOrRegister ( registration ) ;
LwM2mPath resultIds = target ! = null ? new LwM2mPath ( target ) : null ;
if ( ! OBSERVE_READ_ALL . name ( ) . equals ( typeOper . name ( ) ) & & resultIds ! = null & & registration ! = null & & resultIds . getObjectId ( ) > = 0 & & lwM2MClient ! = null ) {
if ( lwM2MClient . isValidObjectVersion ( targetIdVer ) ) {
timeoutInMs = timeoutInMs > 0 ? timeoutInMs : DEFAULT_TIMEOUT ;
ResourceModel resourceModel = null ;
switch ( typeOper ) {
case READ :
request = new ReadRequest ( contentFormat , target ) ;
break ;
case DISCOVER :
request = new DiscoverRequest ( target ) ;
break ;
case OBSERVE :
if ( resultIds . isResource ( ) ) {
request = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) ) ;
} else if ( resultIds . isObjectInstance ( ) ) {
request = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) ) ;
} else if ( resultIds . getObjectId ( ) > = 0 ) {
request = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) ) ;
}
break ;
case OBSERVE_CANCEL :
/ *
lwM2MTransportRequest . sendAllRequest ( lwServer , registration , path , POST_TYPE_OPER_OBSERVE_CANCEL , null , null , null , null , context . getTimeout ( ) ) ;
At server side this will not remove the observation from the observation store , to do it you need to use
{ @code ObservationService # cancelObservation ( ) }
* /
context . getServer ( ) . getObservationService ( ) . cancelObservations ( registration , target ) ;
break ;
case EXECUTE :
resourceModel = lwM2MClient . getResourceModel ( targetIdVer , this . config
. getModelProvider ( ) ) ;
if ( params ! = null & & ! resourceModel . multiple ) {
request = new ExecuteRequest ( target , ( String ) this . converter . convertValue ( params , resourceModel . type , ResourceModel . Type . STRING , resultIds ) ) ;
} else {
request = new ExecuteRequest ( target ) ;
}
break ;
case WRITE_REPLACE :
// Request to write a <b>String Single-Instance Resource</b> using the TLV content format.
resourceModel = lwM2MClient . getResourceModel ( targetIdVer , this . config . getModelProvider ( ) ) ;
if ( contentFormat . equals ( ContentFormat . TLV ) ) {
request = this . getWriteRequestSingleResource ( null , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) , params , resourceModel . type ,
registration , rpcRequest ) ;
}
// Mode.REPLACE && Request to write a <b>String Single-Instance Resource</b> using the given content format (TEXT, TLV, JSON)
else if ( ! contentFormat . equals ( ContentFormat . TLV ) ) {
request = this . getWriteRequestSingleResource ( contentFormat , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) , params , resourceModel . type ,
registration , rpcRequest ) ;
}
break ;
case WRITE_UPDATE :
if ( resultIds . isResource ( ) ) {
/ * *
* send request : path = ' / 3 / 0 ' node = = wM2mObjectInstance
* with params = = "\"resources\" : { 15 : resource : { id : 15 . value : ' + 01 ' . . . } }
* * /
Collection < LwM2mResource > resources = lwM2MClient . getNewOneResourceForInstance (
targetIdVer , params ,
this . config . getModelProvider ( ) ,
this . converter ) ;
request = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resources ) ;
}
/ * *
* params = "{\"id\":0,\"resources\":[{\"id\":14,\"value\":\"+5\"},{\"id\":15,\"value\":\"+9\"}]}"
*
* int rscId = resultIds . getObjectInstanceId ( ) ;
* /
else if ( resultIds . isObjectInstance ( ) ) {
if ( ( ( ConcurrentHashMap ) params ) . size ( ) > 0 ) {
Collection < LwM2mResource > resources = lwM2MClient . getNewManyResourcesForInstance (
targetIdVer , params ,
this . config . getModelProvider ( ) ,
this . converter ) ;
request = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resources ) ;
}
} else if ( resultIds . getObjectId ( ) > = 0 ) {
request = new ObserveRequest ( resultIds . getObjectId ( ) ) ;
}
break ;
case WRITE_ATTRIBUTES :
request = createWriteAttributeRequest ( target , params ) ;
break ;
case DELETE :
request = new DeleteRequest ( target ) ;
break ;
}
DownlinkRequest request = createRequest ( registration , lwM2MClient , typeOper , contentFormat , target ,
targetIdVer , resultIds , params , rpcRequest ) ;
if ( request ! = null ) {
try {
this . sendRequest ( registration , lwM2MClient , request , timeoutInMs , rpcRequest ) ;
} catch ( ClientSleepingException e ) {
DownlinkRequest finalRequest = request ;
long finalTimeoutInMs = timeoutInMs ;
lwM2MClient . getQueuedRequests ( ) . add ( ( ) - > sendRequest ( registration , lwM2MClient , finalRequest , finalTimeoutInMs , rpcRequest ) ) ;
Lwm2mClientRpcRequest finalRpcRequest = rpcRequest ;
lwM2MClient . getQueuedRequests ( ) . add ( ( ) - > sendRequest ( registration , lwM2MClient , finalRequest , finalTimeoutInMs , finalRpcRequest ) ) ;
} catch ( Exception e ) {
log . error ( "[{}] [{}] [{}] Failed to send downlink." , registration . getEndpoint ( ) , targetIdVer , typeOper . name ( ) , e ) ;
}
} else if ( OBSERVE_CANCEL = = typeOper ) {
log . trace ( "[{}], [{}] - [{}] SendRequest" , registration . getEndpoint ( ) , typeOper . name ( ) , targetIdVer ) ;
if ( rpcRequest ! = null ) {
rpcRequest . setInfoMsg ( null ) ;
serviceImpl . sentRpcRequest ( rpcRequest , CONTENT . name ( ) , null , null ) ;
}
else if ( WRITE_UPDATE . name ( ) . equals ( typeOper . name ( ) ) ) {
Lwm2mClientRpcRequest rpcRequestClone = ( Lwm2mClientRpcRequest ) rpcRequest . clone ( ) ;
if ( rpcRequestClone ! = null ) {
String errorMsg = String . format ( "Path %s params is not valid" , targetIdVer ) ;
handler . sentRpcRequest ( rpcRequestClone , BAD_REQUEST . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
rpcRequest = null ;
}
} else {
}
else if ( ! OBSERVE_CANCEL . name ( ) . equals ( typeOper . name ( ) ) ) {
log . error ( "[{}], [{}] - [{}] error SendRequest" , registration . getEndpoint ( ) , typeOper . name ( ) , targetIdVer ) ;
if ( rpcRequest ! = null ) {
ResourceModel resourceModel = lwM2MClient . getResourceModel ( targetIdVer , this . config . getModelProvider ( ) ) ;
String errorMsg = resourceModel = = null ? String . format ( "Path %s not found in object version" , targetIdVer ) : "SendRequest - null" ;
serviceImpl . sentRpcRequest ( rpcRequest , NOT_FOUND . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
handler . sentRpcRequest ( rpcRequest , NOT_FOUND . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
}
}
} else if ( rpcRequest ! = null ) {
String errorMsg = String . format ( "Path %s not found in object version" , targetIdVer ) ;
serviceImpl . sentRpcRequest ( rpcRequest , NOT_FOUND . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
handler . sentRpcRequest ( rpcRequest , NOT_FOUND . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
}
} else if ( OBSERVE_READ_ALL . name ( ) . equals ( typeOper . name ( ) ) | | DISCOVER_All . name ( ) . equals ( typeOper . name ( ) ) ) {
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 ) ;
handler . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
}
} else if ( OBSERVE_READ_ALL . name ( ) . equals ( typeOper . name ( ) ) ) {
Set < Observation > observations = context . getServer ( ) . getObservationService ( ) . getObservations ( registration ) ;
Set < String > observationPaths = observations . stream ( ) . map ( observation - > observation . getPath ( ) . toString ( ) ) . collect ( Collectors . toUnmodifiableSet ( ) ) ;
String msg = String . format ( "%s: type operation %s observation paths - %s" , LOG_LW2M_INFO ,
OBSERVE_READ_ALL . type , observationPaths ) ;
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
log . trace ( "[{}] [{}], [{}]" , typeOper . name ( ) , registration . getEndpoint ( ) , msg ) ;
if ( rpcRequest ! = null ) {
String valueMsg = String . format ( "Observation paths - %s" , observationPaths ) ;
serviceImpl . sentRpcRequest ( rpcRequest , CONTENT . name ( ) , valueMsg , LOG_LW2M_VALUE ) ;
String valueMsg = String . format ( "Paths - %s" , paths ) ;
handler . sentRpcRequest ( rpcRequest , CONTENT . name ( ) , valueMsg , LOG_LW2M_VALUE ) ;
}
} else if ( OBSERVE_CANCEL . name ( ) . equals ( typeOper . name ( ) ) ) {
int observeCancelCnt = context . getServer ( ) . getObservationService ( ) . cancelObservations ( registration ) ;
String observeCancelMsgAll = String . format ( "%s: type operation %s paths: All count: %d" , LOG_LW2M_INFO ,
OBSERVE_CANCEL . name ( ) , observeCancelCnt ) ;
this . afterObserveCancel ( registration , observeCancelCnt , observeCancelMsgAll , rpcRequest ) ;
}
} catch ( Exception e ) {
String msg = String . format ( "%s: type operation %s %s" , LOG_LW2M_ERROR ,
typeOper . name ( ) , e . getMessage ( ) ) ;
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
try {
throw new Exception ( e ) ;
} catch ( Exception exception ) {
exception . printStackTrace ( ) ;
handler . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
if ( rpcRequest ! = null ) {
String errorMsg = String . format ( "Path %s type operation %s %s" , targetIdVer , typeOper . name ( ) , e . getMessage ( ) ) ;
handler . sentRpcRequest ( rpcRequest , NOT_FOUND . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
}
}
}
private DownlinkRequest createRequest ( Registration registration , LwM2mClient lwM2MClient , LwM2mTypeOper typeOper ,
ContentFormat contentFormat , String target , String targetIdVer ,
LwM2mPath resultIds , Object params , Lwm2mClientRpcRequest rpcRequest ) {
DownlinkRequest request = null ;
switch ( typeOper ) {
case READ :
request = new ReadRequest ( contentFormat , target ) ;
break ;
case DISCOVER :
request = new DiscoverRequest ( target ) ;
break ;
case OBSERVE :
String msg = String . format ( "%s: Send Observation %s." , LOG_LW2M_INFO , targetIdVer ) ;
log . warn ( msg ) ;
if ( resultIds . isResource ( ) ) {
Set < Observation > observations = context . getServer ( ) . getObservationService ( ) . getObservations ( registration ) ;
Set < Observation > paths = observations . stream ( ) . filter ( observation - > observation . getPath ( ) . equals ( resultIds ) ) . collect ( Collectors . toSet ( ) ) ;
if ( paths . size ( ) = = 0 ) {
request = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) ) ;
} else {
request = new ReadRequest ( contentFormat , target ) ;
}
} else if ( resultIds . isObjectInstance ( ) ) {
request = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) , resultIds . getObjectInstanceId ( ) ) ;
} else if ( resultIds . getObjectId ( ) > = 0 ) {
request = new ObserveRequest ( contentFormat , resultIds . getObjectId ( ) ) ;
}
break ;
case OBSERVE_CANCEL :
/ *
lwM2MTransportRequest . sendAllRequest ( lwServer , registration , path , POST_TYPE_OPER_OBSERVE_CANCEL , null , null , null , null , context . getTimeout ( ) ) ;
At server side this will not remove the observation from the observation store , to do it you need to use
{ @code ObservationService # cancelObservation ( ) }
* /
int observeCancelCnt = context . getServer ( ) . getObservationService ( ) . cancelObservations ( registration , target ) ;
String observeCancelMsg = String . format ( "%s: type operation %s paths: %s count: %d" , LOG_LW2M_INFO ,
OBSERVE_CANCEL . name ( ) , target , observeCancelCnt ) ;
this . afterObserveCancel ( registration , observeCancelCnt , observeCancelMsg , rpcRequest ) ;
break ;
case EXECUTE :
ResourceModel resourceModelExe = lwM2MClient . getResourceModel ( targetIdVer , this . config . getModelProvider ( ) ) ;
if ( params ! = null & & ! resourceModelExe . multiple ) {
request = new ExecuteRequest ( target , ( String ) this . converter . convertValue ( params , resourceModelExe . type , ResourceModel . Type . STRING , resultIds ) ) ;
} else {
request = new ExecuteRequest ( target ) ;
}
break ;
case WRITE_REPLACE :
/ * *
* Request to write a < b > String Single - Instance Resource < / b > using the TLV content format .
* Type from resourceModel - > STRING , INTEGER , FLOAT , BOOLEAN , OPAQUE , TIME , OBJLNK
* contentFormat - > TLV , TLV , TLV , TLV , OPAQUE , TLV , LINK
* JSON , TEXT ;
* * /
ResourceModel resourceModelWrite = lwM2MClient . getResourceModel ( targetIdVer , this . config . getModelProvider ( ) ) ;
contentFormat = getContentFormatByResourceModelType ( resourceModelWrite , contentFormat ) ;
request = this . getWriteRequestSingleResource ( contentFormat , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resultIds . getResourceId ( ) , params , resourceModelWrite . type ,
registration , rpcRequest ) ;
break ;
case WRITE_UPDATE :
if ( resultIds . isResource ( ) ) {
/ * *
* send request : path = ' / 3 / 0 ' node = = wM2mObjectInstance
* with params = = "\"resources\" : { 15 : resource : { id : 15 . value : ' + 01 ' . . . } }
* * /
Collection < LwM2mResource > resources = lwM2MClient . getNewResourceForInstance (
targetIdVer , params ,
this . config . getModelProvider ( ) ,
this . converter ) ;
contentFormat = getContentFormatByResourceModelType ( lwM2MClient . getResourceModel ( targetIdVer , this . config . getModelProvider ( ) ) ,
contentFormat ) ;
request = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resources ) ;
}
/ * *
* params = "{\"id\":0,\"resources\":[{\"id\":14,\"value\":\"+5\"},{\"id\":15,\"value\":\"+9\"}]}"
* int rscId = resultIds . getObjectInstanceId ( ) ;
* contentFormat – Format of the payload ( TLV or JSON ) .
* /
else if ( resultIds . isObjectInstance ( ) ) {
if ( ( ( ConcurrentHashMap ) params ) . size ( ) > 0 ) {
Collection < LwM2mResource > resources = lwM2MClient . getNewResourcesForInstance (
targetIdVer , params ,
this . config . getModelProvider ( ) ,
this . converter ) ;
if ( resources . size ( ) > 0 ) {
contentFormat = contentFormat . equals ( ContentFormat . JSON ) ? contentFormat : ContentFormat . TLV ;
request = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resources ) ;
}
}
} else if ( resultIds . getObjectId ( ) > = 0 ) {
request = new ObserveRequest ( resultIds . getObjectId ( ) ) ;
}
break ;
case WRITE_ATTRIBUTES :
request = createWriteAttributeRequest ( target , params ) ;
break ;
case DELETE :
request = new DeleteRequest ( target ) ;
break ;
}
return request ;
}
/ * *
* @param registration -
* @param request -
@ -273,52 +316,66 @@ public class LwM2mTransportRequest {
private void sendRequest ( Registration registration , LwM2mClient lwM2MClient , DownlinkRequest request ,
long timeoutInMs , Lwm2mClientRpcRequest rpcRequest ) {
context . getServer ( ) . send ( registration , request , timeoutInMs , ( ResponseCallback < ? > ) response - > {
if ( ! lwM2MClient . isInit ( ) ) {
lwM2MClient . initReadValue ( this . serviceImpl , convertPathFromObjectIdToIdVer ( request . getPath ( ) . toString ( ) , registration ) ) ;
lwM2MClient . initReadValue ( this . handler , convertPathFromObjectIdToIdVer ( request . getPath ( ) . toString ( ) , registration ) ) ;
}
if ( CoAP . ResponseCode . isSuccess ( ( ( Response ) response . getCoapResponse ( ) ) . getCode ( ) ) ) {
this . handleResponse ( registration , 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 ( ) ) ;
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
handler . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
log . error ( "[{}] [{}], [{}] - [{}] [{}] error SendRequest" , request . getClass ( ) . getName ( ) . toString ( ) , registration . getEndpoint ( ) ,
( ( Response ) response . getCoapResponse ( ) ) . getCode ( ) , response . getCode ( ) , request . getPath ( ) . toString ( ) ) ;
if ( ! lwM2MClient . isInit ( ) ) {
lwM2MClient . initReadValue ( this . serviceImpl , convertPathFromObjectIdToIdVer ( request . getPath ( ) . toString ( ) , registration ) ) ;
lwM2MClient . initReadValue ( this . handler , convertPathFromObjectIdToIdVer ( request . getPath ( ) . toString ( ) , registration ) ) ;
}
/** Not Found */
if ( rpcRequest ! = null ) {
serviceImpl . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , response . getErrorMessage ( ) , LOG_LW2M_ERROR ) ;
handler . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , response . getErrorMessage ( ) , LOG_LW2M_ERROR ) ;
}
/ * Not Found
set setClient_fw_version = empty
* /
if ( FR_PATH_RESOURCE_VER_ID . equals ( request . getPath ( ) . toString ( ) ) & & lwM2MClient . isUpdateFw ( ) ) {
lwM2MClient . setUpdateFw ( false ) ;
lwM2MClient . getFrUpdate ( ) . setClientFwVersion ( "" ) ;
log . warn ( "updateFirmwareClient1" ) ;
serviceImpl . updateFirmwareClient ( lwM2MClient ) ;
/ * * Not Found
set setClient_fw_info . . . = empty
* * /
if ( lwM2MClient . getFwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getFwUpdate ( ) . initReadValue ( handler , this , request . getPath ( ) . toString ( ) ) ;
}
if ( lwM2MClient . getSwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getSwUpdate ( ) . initReadValue ( handler , this , request . getPath ( ) . toString ( ) ) ;
}
if ( request . getPath ( ) . toString ( ) . equals ( FW_PACKAGE_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 - > {
/ * version = = null
set setClient_fw_version = empty
* /
if ( FR_PATH_RESOURCE_VER_ID . equals ( request . getPath ( ) . toString ( ) ) & & lwM2MClient . isUpdateFw ( ) ) {
lwM2MClient . setUpdateFw ( false ) ;
lwM2MClient . getFrUpdate ( ) . setClientFwVersion ( "" ) ;
log . warn ( "updateFirmwareClient2" ) ;
serviceImpl . updateFirmwareClient ( lwM2MClient ) ;
/ * * version = = null
set setClient_fw_info . . . = empty
* * /
if ( lwM2MClient . getFwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getFwUpdate ( ) . initReadValue ( handler , this , request . getPath ( ) . toString ( ) ) ;
}
if ( lwM2MClient . getSwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getSwUpdate ( ) . initReadValue ( handler , this , request . getPath ( ) . toString ( ) ) ;
}
if ( request . getPath ( ) . toString ( ) . equals ( FW_PACKAGE_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 ( ! lwM2MClient . isInit ( ) ) {
lwM2MClient . initReadValue ( this . serviceImpl , convertPathFromObjectIdToIdVer ( request . getPath ( ) . toString ( ) , registration ) ) ;
lwM2MClient . 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 ( ) ) ;
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
handler . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
log . error ( "[{}] [{}] - [{}] error SendRequest" , request . getClass ( ) . getName ( ) . toString ( ) , request . getPath ( ) . toString ( ) , e . toString ( ) ) ;
if ( rpcRequest ! = null ) {
serviceImpl . sentRpcRequest ( rpcRequest , CoAP . CodeClass . ERROR_RESPONSE . name ( ) , e . getMessage ( ) , LOG_LW2M_ERROR ) ;
handler . sentRpcRequest ( rpcRequest , CoAP . CodeClass . ERROR_RESPONSE . name ( ) , e . getMessage ( ) , LOG_LW2M_ERROR ) ;
}
} ) ;
}
@ -360,11 +417,11 @@ public class LwM2mTransportRequest {
String patn = "/" + objectId + "/" + instanceId + "/" + resourceId ;
String msg = String . format ( LOG_LW2M_ERROR + ": NumberFormatException: Resource path - %s type - %s value - %s msg error - %s SendRequest to Client" ,
patn , type , value , e . toString ( ) ) ;
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
handler . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
log . error ( "Path: [{}] type: [{}] value: [{}] errorMsg: [{}]]" , patn , type , value , e . toString ( ) ) ;
if ( rpcRequest ! = null ) {
String errorMsg = String . format ( "NumberFormatException: Resource path - %s type - %s value - %s" , patn , type , value ) ;
serviceImpl . sentRpcRequest ( rpcRequest , BAD_REQUEST . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
handler . sentRpcRequest ( rpcRequest , BAD_REQUEST . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
}
return null ;
}
@ -372,7 +429,7 @@ public class LwM2mTransportRequest {
private void handleResponse ( Registration registration , final String path , LwM2mResponse response ,
DownlinkRequest request , Lwm2mClientRpcRequest rpcRequest ) {
executorResponse . submit ( ( ) - > {
responseRequ estE xecutor. submit ( ( ) - > {
try {
this . sendResponse ( registration , path , response , request , rpcRequest ) ;
} catch ( Exception e ) {
@ -391,39 +448,37 @@ public class LwM2mTransportRequest {
private void sendResponse ( Registration registration , String path , LwM2mResponse response ,
DownlinkRequest request , Lwm2mClientRpcRequest rpcRequest ) {
String pathIdVer = convertPathFromObjectIdToIdVer ( path , registration ) ;
String msgLog = "" ;
if ( response instanceof ReadResponse ) {
serviceImpl . onUpdateValueAfterReadResponse ( registration , pathIdVer , ( ReadResponse ) response , rpcRequest ) ;
} else if ( response instanceof CancelObservationResponse ) {
log . info ( "[{}] Path [{}] CancelObservationResponse 3_Send" , pathIdVer , response ) ;
handler . onUpdateValueAfterReadResponse ( registration , pathIdVer , ( ReadResponse ) response , rpcRequest ) ;
} else if ( response instanceof DeleteResponse ) {
log . info ( "[{}] Path [{}] DeleteResponse 5_Send" , pathIdVer , response ) ;
log . warn ( "[{}] Path [{}] DeleteResponse 5_Send" , pathIdVer , response ) ;
} else if ( response instanceof DiscoverResponse ) {
log . info ( "[{}] [{}] - [{}] [{}] Discovery value: [{}]" , registration . getEndpoint ( ) ,
( ( Response ) response . getCoapResponse ( ) ) . getCode ( ) , response . getCode ( ) ,
request . getPath ( ) . toString ( ) , ( ( DiscoverResponse ) response ) . getObjectLinks ( ) ) ;
String discoverValue = Link . serialize ( ( ( DiscoverResponse ) response ) . getObjectLinks ( ) ) ;
msgLog = String . format ( "%s: type operation: %s path: %s value: %s" ,
LOG_LW2M_INFO , DISCOVER . name ( ) , request . getPath ( ) . toString ( ) , discoverValue ) ;
handler . sendLogsToThingsboard ( msgLog , registration . getId ( ) ) ;
log . warn ( "DiscoverResponse: [{}]" , ( DiscoverResponse ) response ) ;
if ( rpcRequest ! = null ) {
String discoveryMsg = String . format ( "%s" ,
Arrays . stream ( ( ( DiscoverResponse ) response ) . getObjectLinks ( ) ) . collect ( Collectors . toSet ( ) ) ) ;
serviceImpl . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , discoveryMsg , LOG_LW2M_VALUE ) ;
handler . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , discoverValue , LOG_LW2M_VALUE ) ;
}
} else if ( response instanceof ExecuteResponse ) {
log . info ( "[{}] Path [{}] ExecuteResponse 7_Send" , pathIdVer , response ) ;
log . warn ( "[{}] Path [{}] ExecuteResponse 7_Send" , pathIdVer , response ) ;
} else if ( response instanceof WriteAttributesResponse ) {
log . info ( "[{}] Path [{}] WriteAttributesResponse 8_Send" , pathIdVer , response ) ;
log . warn ( "[{}] Path [{}] WriteAttributesResponse 8_Send" , pathIdVer , response ) ;
} else if ( response instanceof WriteResponse ) {
log . info ( "[{}] Path [{}] WriteResponse 9_Send" , pathIdVer , response ) ;
log . warn ( "[{}] Path [{}] WriteResponse 9_Send" , pathIdVer , response ) ;
this . infoWriteResponse ( registration , response , request ) ;
serviceImpl . onWriteResponseOk ( registration , pathIdVer , ( WriteRequest ) request ) ;
handler . onWriteResponseOk ( registration , pathIdVer , ( WriteRequest ) request ) ;
}
if ( rpcRequest ! = null ) {
if ( response instanceof ExecuteResponse
| | response instanceof WriteAttributesResponse
| | response instanceof DeleteResponse ) {
rpcRequest . setInfoMsg ( null ) ;
serviceImpl . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , null , null ) ;
handler . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , null , null ) ;
} else if ( response instanceof WriteResponse ) {
serviceImpl . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , null , LOG_LW2M_INFO ) ;
handler . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , null , LOG_LW2M_INFO ) ;
}
}
}
@ -447,21 +502,73 @@ public class LwM2mTransportRequest {
Math . min ( valueLength , config . getLogMaxLength ( ) ) ) ) ;
}
value = valueLength > config . getLogMaxLength ( ) ? value + "..." : value ;
msg = String . format ( "%s: Update finished successfully: Lwm2m code - %d Resource path - %s length - %s value - %s" ,
msg = String . format ( "%s: Update finished successfully: Lwm2m code - %d Resource path: %s length: %s value: %s" ,
LOG_LW2M_INFO , response . getCode ( ) . getCode ( ) , request . getPath ( ) . toString ( ) , valueLength , value ) ;
} else {
value = this . converter . convertValue ( singleResource . getValue ( ) ,
singleResource . getType ( ) , ResourceModel . Type . STRING , request . getPath ( ) ) ;
msg = String . format ( "%s: Update finished successfully: Lwm2m code - %d Resource path - %s value - %s" ,
msg = String . format ( "%s: Update finished successfully. Lwm2m code: %d Resource path: %s value: %s" ,
LOG_LW2M_INFO , response . getCode ( ) . getCode ( ) , request . getPath ( ) . toString ( ) , value ) ;
}
if ( msg ! = null ) {
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
log . warn ( "[{}] [{}] [{}] - [{}] [{}] Update finished successfully: [{}]" , request . getClass ( ) . getName ( ) , registration . getEndpoint ( ) ,
( ( Response ) response . getCoapResponse ( ) ) . getCode ( ) , response . getCode ( ) , request . getPath ( ) . toString ( ) , value ) ;
handler . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
if ( request . getPath ( ) . toString ( ) . equals ( FW_PACKAGE_ID ) | | request . getPath ( ) . toString ( ) . equals ( SW_PACKAGE_ID ) ) {
this . afterWriteSuccessFwSwUpdate ( registration , request ) ;
}
}
} catch ( Exception e ) {
log . trace ( "Fail convert value from request to string. " , e ) ;
}
}
/ * *
* After finish operation FwSwUpdate Write ( success ) :
* fw_state / sw_state = DOWNLOADED
* send operation Execute
* /
private void afterWriteSuccessFwSwUpdate ( Registration registration , DownlinkRequest request ) {
LwM2mClient lwM2MClient = this . lwM2mClientContext . getClientByRegistrationId ( registration . getId ( ) ) ;
if ( request . getPath ( ) . toString ( ) . equals ( FW_PACKAGE_ID ) & & lwM2MClient . getFwUpdate ( ) ! = null ) {
lwM2MClient . getFwUpdate ( ) . setStateUpdate ( DOWNLOADED . name ( ) ) ;
lwM2MClient . getFwUpdate ( ) . sendLogs ( handler , WRITE_REPLACE . name ( ) , LOG_LW2M_INFO , null ) ;
}
if ( request . getPath ( ) . toString ( ) . equals ( SW_PACKAGE_ID ) & & lwM2MClient . getSwUpdate ( ) ! = null ) {
lwM2MClient . getSwUpdate ( ) . setStateUpdate ( DOWNLOADED . name ( ) ) ;
lwM2MClient . getSwUpdate ( ) . sendLogs ( handler , WRITE_REPLACE . name ( ) , LOG_LW2M_INFO , null ) ;
}
}
/ * *
* After finish operation FwSwUpdate Write ( error ) : fw_state = FAILED
* /
private void afterWriteFwSWUpdateError ( Registration registration , DownlinkRequest request , String msgError ) {
LwM2mClient lwM2MClient = this . lwM2mClientContext . getClientByRegistrationId ( registration . getId ( ) ) ;
if ( request . getPath ( ) . toString ( ) . equals ( FW_PACKAGE_ID ) & & lwM2MClient . getFwUpdate ( ) ! = null ) {
lwM2MClient . getFwUpdate ( ) . setStateUpdate ( FAILED . name ( ) ) ;
lwM2MClient . getFwUpdate ( ) . sendLogs ( handler , WRITE_REPLACE . name ( ) , LOG_LW2M_ERROR , msgError ) ;
}
if ( request . getPath ( ) . toString ( ) . equals ( SW_PACKAGE_ID ) & & lwM2MClient . getSwUpdate ( ) ! = null ) {
lwM2MClient . getSwUpdate ( ) . setStateUpdate ( FAILED . name ( ) ) ;
lwM2MClient . getSwUpdate ( ) . sendLogs ( handler , WRITE_REPLACE . name ( ) , LOG_LW2M_ERROR , msgError ) ;
}
}
private void afterExecuteFwSwUpdateError ( Registration registration , DownlinkRequest request , String msgError ) {
LwM2mClient lwM2MClient = this . lwM2mClientContext . getClientByRegistrationId ( registration . getId ( ) ) ;
if ( request . getPath ( ) . toString ( ) . equals ( FW_UPDATE_ID ) & & lwM2MClient . getFwUpdate ( ) ! = null ) {
lwM2MClient . getFwUpdate ( ) . sendLogs ( handler , EXECUTE . name ( ) , LOG_LW2M_ERROR , msgError ) ;
}
if ( request . getPath ( ) . toString ( ) . equals ( SW_INSTALL_ID ) & & lwM2MClient . getSwUpdate ( ) ! = null ) {
lwM2MClient . getSwUpdate ( ) . sendLogs ( handler , EXECUTE . name ( ) , LOG_LW2M_ERROR , msgError ) ;
}
}
private void afterObserveCancel ( Registration registration , int observeCancelCnt , String observeCancelMsg , Lwm2mClientRpcRequest rpcRequest ) {
handler . sendLogsToThingsboard ( observeCancelMsg , registration . getId ( ) ) ;
log . warn ( "[{}]" , observeCancelMsg ) ;
if ( rpcRequest ! = null ) {
rpcRequest . setInfoMsg ( String . format ( "Count: %d" , observeCancelCnt ) ) ;
handler . sentRpcRequest ( rpcRequest , CONTENT . name ( ) , null , LOG_LW2M_INFO ) ;
}
}
}