@ -36,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 ;
@ -49,7 +48,6 @@ import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.core.util.NamedThreadFactory ;
import org.eclipse.leshan.server.registration.Registration ;
import org.springframework.stereotype.Service ;
import org.thingsboard.server.common.data.firmware.FirmwareUpdateStatus ;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent ;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig ;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient ;
@ -70,16 +68,25 @@ 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.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 ;
@ -90,7 +97,7 @@ 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 ;
@ -102,8 +109,8 @@ public class LwM2mTransportRequest {
@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 ) ) ) ;
}
/ * *
@ -119,118 +126,14 @@ 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 ( ) ) {
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 ( ) }
* /
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 ) ;
if ( resources . size ( ) > 0 ) {
request = new WriteRequest ( WriteRequest . Mode . UPDATE , contentFormat , resultIds . getObjectId ( ) ,
resultIds . getObjectInstanceId ( ) , resources ) ;
} else {
Lwm2mClientRpcRequest rpcRequestClone = ( Lwm2mClientRpcRequest ) rpcRequest . clone ( ) ;
if ( rpcRequestClone ! = null ) {
String errorMsg = String . format ( "Path %s params is not valid" , targetIdVer ) ;
serviceImpl . sentRpcRequest ( rpcRequestClone , BAD_REQUEST . getName ( ) , errorMsg , LOG_LW2M_ERROR ) ;
rpcRequest = null ;
}
}
}
} 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 ) ;
@ -242,15 +145,19 @@ public class LwM2mTransportRequest {
} 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 ) ;
serviceImpl . 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 ) ;
}
@ -265,30 +172,140 @@ public class LwM2mTransportRequest {
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 - > link . toString ( ) ) . collect ( Collectors . toUnmodifiableSet ( ) ) ;
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 ) ;
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
}
String msg = String . format ( "%s: type operation %s paths - %s" , LOG_LW2M_INFO ,
OBSERVE_READ_ALL . type , paths ) ;
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
log . warn ( "[{}] [{}], [{}]" , typeOper . name ( ) , registration . getEndpoint ( ) , msg ) ;
if ( rpcRequest ! = null ) {
String valueMsg = String . format ( "Paths - %s" , paths ) ;
serviceImpl . 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 ( ) ;
if ( rpcRequest ! = null ) {
String errorMsg = String . format ( "Path %s type operation %s %s" , targetIdVer , typeOper . name ( ) , e . getMessage ( ) ) ;
serviceImpl . 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 -
@ -314,6 +331,7 @@ public class LwM2mTransportRequest {
if ( ! lwM2MClient . isInit ( ) ) {
lwM2MClient . initReadValue ( this . serviceImpl , convertPathFromObjectIdToIdVer ( request . getPath ( ) . toString ( ) , registration ) ) ;
}
/** Not Found */
if ( rpcRequest ! = null ) {
serviceImpl . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , response . getErrorMessage ( ) , LOG_LW2M_ERROR ) ;
}
@ -322,7 +340,15 @@ public class LwM2mTransportRequest {
* * /
if ( lwM2MClient . getFwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getFwUpdate ( ) . initReadValue ( serviceImpl , request . getPath ( ) . toString ( ) ) ;
log . warn ( "updateFirmwareClient1" ) ;
}
if ( lwM2MClient . getSwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getSwUpdate ( ) . initReadValue ( serviceImpl , 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 - > {
@ -331,7 +357,15 @@ public class LwM2mTransportRequest {
* * /
if ( lwM2MClient . getFwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getFwUpdate ( ) . initReadValue ( serviceImpl , request . getPath ( ) . toString ( ) ) ;
log . warn ( "updateFirmwareClient2" ) ;
}
if ( lwM2MClient . getSwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getSwUpdate ( ) . initReadValue ( serviceImpl , 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 ) ) ;
@ -395,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 ) {
@ -414,28 +448,26 @@ 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 ) ;
} 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 ) ;
serviceImpl . 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 ) ;
serviceImpl . 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 ) ;
}
@ -480,9 +512,8 @@ public class LwM2mTransportRequest {
}
if ( msg ! = null ) {
serviceImpl . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
log . warn ( msg ) ;
if ( request . getPath ( ) . toString ( ) . equals ( FW_PACKAGE_ID ) | | request . getPath ( ) . toString ( ) . equals ( SW_PACKAGE_ID ) ) {
this . execute FwSwUpdate( registration , request ) ;
this . afterWriteSuccess FwSwUpdate( registration , request ) ;
}
}
} catch ( Exception e ) {
@ -490,15 +521,54 @@ public class LwM2mTransportRequest {
}
}
private void executeFwSwUpdate ( Registration registration , DownlinkRequest request ) {
LwM2mClient lwM2mClient = this . lwM2mClientContext . getClientByRegistrationId ( registration . getId ( ) ) ;
if ( request . getPath ( ) . toString ( ) . equals ( FW_PACKAGE_ID )
& & FirmwareUpdateStatus . DOWNLOADING . name ( ) . equals ( lwM2mClient . getFwUpdate ( ) . getStateUpdate ( ) ) ) {
lwM2mClient . getFwUpdate ( ) . sendReadInfoForWrite ( ) ;
/ * *
* 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 ( 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 ( 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 ( 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 ( WRITE_REPLACE . name ( ) , LOG_LW2M_ERROR , msgError ) ;
}
if ( request . getPath ( ) . toString ( ) . equals ( SW_PACKAGE_ID )
& & FirmwareUpdateStatus . DOWNLOADING . name ( ) . equals ( lwM2mClient . getSwUpdate ( ) . getStateUpdate ( ) ) ) {
lwM2mClient . getSwUpdate ( ) . sendReadInfoForWrite ( ) ;
}
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 ( EXECUTE . name ( ) , LOG_LW2M_ERROR , msgError ) ;
}
if ( request . getPath ( ) . toString ( ) . equals ( SW_INSTALL_ID ) & & lwM2MClient . getSwUpdate ( ) ! = null ) {
lwM2MClient . getSwUpdate ( ) . sendLogs ( EXECUTE . name ( ) , LOG_LW2M_ERROR , msgError ) ;
}
}
private void afterObserveCancel ( Registration registration , int observeCancelCnt , String observeCancelMsg , Lwm2mClientRpcRequest rpcRequest ) {
serviceImpl . sendLogsToThingsboard ( observeCancelMsg , registration . getId ( ) ) ;
log . warn ( "[{}]" , observeCancelMsg ) ;
if ( rpcRequest ! = null ) {
rpcRequest . setInfoMsg ( String . format ( "Count: %d" , observeCancelCnt ) ) ;
serviceImpl . sentRpcRequest ( rpcRequest , CONTENT . name ( ) , null , LOG_LW2M_INFO ) ;
}
}
}