@ -23,6 +23,7 @@ import com.google.gson.JsonElement;
import com.google.gson.JsonObject ;
import com.google.gson.reflect.TypeToken ;
import lombok.extern.slf4j.Slf4j ;
import org.apache.commons.lang3.StringUtils ;
import org.eclipse.leshan.core.model.ObjectModel ;
import org.eclipse.leshan.core.model.ResourceModel ;
import org.eclipse.leshan.core.node.LwM2mObject ;
@ -38,13 +39,13 @@ import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service ;
import org.thingsboard.common.util.JacksonUtil ;
import org.thingsboard.common.util.ThingsBoardExecutors ;
import org.thingsboard.server.cache.firmware.Firmwar eDataCache ;
import org.thingsboard.server.cache.ota.OtaPackag eDataCache ;
import org.thingsboard.server.common.data.Device ;
import org.thingsboard.server.common.data.DeviceProfile ;
import org.thingsboard.server.common.data.firmware.Firmwar eKey ;
import org.thingsboard.server.common.data.firmware.Firmwar eType ;
import org.thingsboard.server.common.data.firmware.Firmwar eUtil ;
import org.thingsboard.server.common.data.id.Firmwar eId ;
import org.thingsboard.server.common.data.ota.OtaPackag eKey ;
import org.thingsboard.server.common.data.ota.OtaPackag eType ;
import org.thingsboard.server.common.data.ota.OtaPackag eUtil ;
import org.thingsboard.server.common.data.id.OtaPackag eId ;
import org.thingsboard.server.common.transport.TransportService ;
import org.thingsboard.server.common.transport.TransportServiceCallback ;
import org.thingsboard.server.common.transport.adaptor.AdaptorException ;
@ -61,6 +62,7 @@ import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext ;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientProfile ;
import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcRequest ;
import org.thingsboard.server.transport.lwm2m.server.client.ResourceValue ;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAddKeyValueProto ;
import org.thingsboard.server.transport.lwm2m.server.client.ResultsAnalyzerParameters ;
import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore ;
@ -85,8 +87,8 @@ import java.util.stream.Collectors;
import static org.eclipse.californium.core.coap.CoAP.ResponseCode.BAD_REQUEST ;
import static org.eclipse.leshan.core.attributes.Attribute.OBJECT_VERSION ;
import static org.thingsboard.server.common.data.firmware.Firmwar eUpdateStatus.DOWNLOADED ;
import static org.thingsboard.server.common.data.firmware.Firmwar eUpdateStatus.UPDATING ;
import static org.thingsboard.server.common.data.ota.OtaPackag eUpdateStatus.DOWNLOADED ;
import static org.thingsboard.server.common.data.ota.OtaPackag eUpdateStatus.UPDATING ;
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getValueFromKvProto ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.CLIENT_NOT_AUTHORIZED ;
@ -99,21 +101,21 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.L
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LW2M_VALUE ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LWM2M_STRATEGY_2 ;
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 ;
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.LwM2mTypeOper.OBSERVE_CANCEL _ALL ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.READ ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LwM2mTypeOper.WRITE_ATTRIBUTES ;
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.SW_ID ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SW_RESULT_ID ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertJsonArrayToSet ;
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.getAckCallback ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.isFwSwWords ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.setValidTypeOper ;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validateObjectVerFromKey ;
@ -125,32 +127,33 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
private ExecutorService registrationExecutor ;
private ExecutorService updateRegistrationExecutor ;
private ExecutorService unRegistrationExecutor ;
private LwM2mValueConverterImpl converter ;
public LwM2mValueConverterImpl converter ;
private final TransportService transportService ;
private final LwM2mTransportContext context ;
public final LwM2MTransportServerConfig config ;
public final FirmwareDataCache firmwar eDataCache;
public final OtaPackageDataCache otaPackag eDataCache;
public final LwM2mTransportServerHelper helper ;
private final LwM2MJsonAdaptor adaptor ;
private final LwM2mClientContext clientContext ;
private final LwM2mTransportRequest lwM2mTransportRequest ;
private final TbLwM2MDtlsSessionStore sessionStore ;
public final LwM2mClientContext clientContext ;
public final LwM2mTransportRequest lwM2mTransportRequest ;
private final Map < UUID , Long > rpcSubscriptions ;
public DefaultLwM2MTransportMsgHandler ( TransportService transportService , LwM2MTransportServerConfig config , LwM2mTransportServerHelper helper ,
LwM2mClientContext clientContext ,
@Lazy LwM2mTransportRequest lwM2mTransportRequest ,
FirmwareDataCache firmwar eDataCache,
OtaPackageDataCache otaPackag eDataCache,
LwM2mTransportContext context , LwM2MJsonAdaptor adaptor , TbLwM2MDtlsSessionStore sessionStore ) {
this . transportService = transportService ;
this . config = config ;
this . helper = helper ;
this . clientContext = clientContext ;
this . lwM2mTransportRequest = lwM2mTransportRequest ;
this . firmwareDataCache = firmwar eDataCache;
this . otaPackageDataCache = otaPackag eDataCache;
this . context = context ;
this . adaptor = adaptor ;
this . rpcSubscriptions = new ConcurrentHashMap < > ( ) ;
this . sessionStore = sessionStore ;
}
@ -193,8 +196,8 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
. setSubscribeToRPC ( TransportProtos . SubscribeToRPCMsg . newBuilder ( ) . build ( ) )
. build ( ) ;
transportService . process ( msg , null ) ;
this . getInfoFirmwareUpdate ( lwM2MClient ) ;
this . getInfoSoftwareUpdate ( lwM2MClient ) ;
this . getInfoFirmwareUpdate ( lwM2MClient , null ) ;
this . getInfoSoftwareUpdate ( lwM2MClient , null ) ;
this . initLwM2mFromClientValue ( registration , lwM2MClient ) ;
this . sendLogsToThingsboard ( LOG_LW2M_INFO + ": Client create after Registration" , registration . getId ( ) ) ;
} else {
@ -241,16 +244,13 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
/ * *
* @param registration - Registration LwM2M Client
* @param observations - All paths observations before unReg
* ! ! ! Warn : if have not finishing unReg , then this operation will be finished on next Client ` s connect
* @param observations - ! ! ! Warn : if have not finishing unReg , then this operation will be finished on next Client ` s connect
* /
public void unReg ( Registration registration , Collection < Observation > observations ) {
unRegistrationExecutor . submit ( ( ) - > {
try {
this . setCancelObservationsAll ( registration ) ;
this . sendLogsToThingsboard ( LOG_LW2M_INFO + ": Client unRegistration" , registration . getId ( ) ) ;
this . closeClientSession ( registration ) ;
;
} catch ( Throwable t ) {
log . error ( "[{}] endpoint [{}] error Unable un registration." , registration . getEndpoint ( ) , t ) ;
this . sendLogsToThingsboard ( LOG_LW2M_ERROR + String . format ( ": Client Unable un Registration, %s" , t . getMessage ( ) ) , registration . getId ( ) ) ;
@ -285,12 +285,8 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
@Override
public void setCancelObservationsAll ( Registration registration ) {
if ( registration ! = null ) {
lwM2mTransportRequest . sendAllRequest ( registration , null , OBSERVE_CANCEL ,
this . lwM2mTransportRequest . sendAllRequest ( registration , null , OBSERVE_CANCEL_AL L ,
null , null , this . config . getTimeout ( ) , null ) ;
// Set<Observation> observations = context.getServer().getObservationService().getObservations(registration);
// observations.forEach(observation -> lwM2mTransportRequest.sendAllRequest(registration,
// convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), OBSERVE_CANCEL,
// null, null, this.config.getTimeout(), null));
}
}
@ -338,7 +334,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
READ , pathIdVer , value ) ;
this . sendLogsToThingsboard ( msg , registration . getId ( ) ) ;
rpcRequest . setValueMsg ( String . format ( "%s" , value ) ) ;
this . sentRpcRequest ( rpcRequest , response . getCode ( ) . getName ( ) , ( String ) value , LOG_LW2M_VALUE ) ;
this . sentRpcResponse ( rpcRequest , response . getCode ( ) . getName ( ) , ( String ) value , LOG_LW2M_VALUE ) ;
}
/ * *
@ -354,22 +350,23 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
* /
@Override
public void onAttributeUpdate ( AttributeUpdateNotificationMsg msg , TransportProtos . SessionInfoProto sessionInfo ) {
LwM2mClient lwM2MClient = clientContext . getClient ( new UUID ( sessionInfo . getSessionIdMSB ( ) , sessionInfo . getSessionIdLSB ( ) ) ) ;
if ( msg . getSharedUpdatedCount ( ) > 0 ) {
LwM2mClient lwM2MClient = clientContext . getClient ( sessionInfo ) ;
if ( msg . getSharedUpdatedCount ( ) > 0 & & lwM2MClient ! = null ) {
log . warn ( "2) OnAttributeUpdate, SharedUpdatedList() [{}]" , msg . getSharedUpdatedList ( ) ) ;
msg . getSharedUpdatedList ( ) . forEach ( tsKvProto - > {
String pathName = tsKvProto . getKv ( ) . getKey ( ) ;
String pathIdVer = this . getPresentPathIntoProfile ( sessionInfo , pathName ) ;
Object valueNew = getValueFromKvProto ( tsKvProto . getKv ( ) ) ;
if ( ( Firmwar eUtil. getAttributeKey ( Firmwar eType. FIRMWARE , Firmwar eKey. VERSION ) . equals ( pathName )
if ( ( OtaPackag eUtil. getAttributeKey ( OtaPackag eType. FIRMWARE , OtaPackag eKey. VERSION ) . equals ( pathName )
& & ( ! valueNew . equals ( lwM2MClient . getFwUpdate ( ) . getCurrentVersion ( ) ) ) )
| | ( Firmwar eUtil. getAttributeKey ( Firmwar eType. FIRMWARE , Firmwar eKey. TITLE ) . equals ( pathName )
| | ( OtaPackag eUtil. getAttributeKey ( OtaPackag eType. FIRMWARE , OtaPackag eKey. TITLE ) . equals ( pathName )
& & ( ! valueNew . equals ( lwM2MClient . getFwUpdate ( ) . getCurrentTitle ( ) ) ) ) ) {
this . getInfoFirmwareUpdate ( lwM2MClient ) ;
} else if ( ( Firmwar eUtil. getAttributeKey ( Firmwar eType. SOFTWARE , Firmwar eKey. VERSION ) . equals ( pathName )
this . getInfoFirmwareUpdate ( lwM2MClient , null ) ;
} else if ( ( OtaPackag eUtil. getAttributeKey ( OtaPackag eType. SOFTWARE , OtaPackag eKey. VERSION ) . equals ( pathName )
& & ( ! valueNew . equals ( lwM2MClient . getSwUpdate ( ) . getCurrentVersion ( ) ) ) )
| | ( Firmwar eUtil. getAttributeKey ( Firmwar eType. SOFTWARE , Firmwar eKey. TITLE ) . equals ( pathName )
| | ( OtaPackag eUtil. getAttributeKey ( OtaPackag eType. SOFTWARE , OtaPackag eKey. TITLE ) . equals ( pathName )
& & ( ! valueNew . equals ( lwM2MClient . getSwUpdate ( ) . getCurrentTitle ( ) ) ) ) ) {
this . getInfoSoftwareUpdate ( lwM2MClient ) ;
this . getInfoSoftwareUpdate ( lwM2MClient , null ) ;
}
if ( pathIdVer ! = null ) {
ResourceModel resourceModel = lwM2MClient . getResourceModel ( pathIdVer , this . config
@ -382,7 +379,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
LOG_LW2M_ERROR , pathIdVer , valueNew ) ;
this . sendLogsToThingsboard ( logMsg , lwM2MClient . getRegistration ( ) . getId ( ) ) ;
}
} else {
} else if ( ! isFwSwWords ( pathName ) ) {
log . error ( "Resource name name - [{}] value - [{}] is not present as attribute/telemetry in profile and cannot be updated" , pathName , valueNew ) ;
String logMsg = String . format ( "%s: attributeUpdate: attribute name - %s value - %s is not present as attribute in profile and cannot be updated" ,
LOG_LW2M_ERROR , pathName , valueNew ) ;
@ -390,16 +387,19 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
}
} ) ;
} else if ( msg . getSharedDeletedCount ( ) > 0 ) {
} else if ( msg . getSharedDeletedCount ( ) > 0 & & lwM2MClient ! = null ) {
msg . getSharedUpdatedList ( ) . forEach ( tsKvProto - > {
String pathName = tsKvProto . getKv ( ) . getKey ( ) ;
Object valueNew = getValueFromKvProto ( tsKvProto . getKv ( ) ) ;
if ( Firmwar eUtil. getAttributeKey ( Firmwar eType. FIRMWARE , Firmwar eKey. VERSION ) . equals ( pathName ) & & ! valueNew . equals ( lwM2MClient . getFwUpdate ( ) . getCurrentVersion ( ) ) ) {
if ( OtaPackag eUtil. getAttributeKey ( OtaPackag eType. FIRMWARE , OtaPackag eKey. VERSION ) . equals ( pathName ) & & ! valueNew . equals ( lwM2MClient . getFwUpdate ( ) . getCurrentVersion ( ) ) ) {
lwM2MClient . getFwUpdate ( ) . setCurrentVersion ( ( String ) valueNew ) ;
}
} ) ;
log . info ( "[{}] delete [{}] onAttributeUpdate" , msg . getSharedDeletedList ( ) , sessionInfo ) ;
}
else if ( lwM2MClient = = null ) {
log . error ( "OnAttributeUpdate, lwM2MClient is null" ) ;
}
}
/ * *
@ -438,124 +438,52 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
clientContext . getLwM2mClients ( ) . forEach ( e - > e . deleteResources ( pathIdVer , this . config . getModelProvider ( ) ) ) ;
}
/ * *
* # 1 del from rpcSubscriptions by timeout
* # 2 if not present in rpcSubscriptions by requestId : create new Lwm2mClientRpcRequest , after success - add requestId , timeout
* /
@Override
public void onToDeviceRpcRequest ( TransportProtos . ToDeviceRpcRequestMsg toDeviceRpcRequestMsg , SessionInfoProto sessionInfo ) {
log . warn ( "4) RPC-OK finish to [{}]" , toDeviceRpcRequestMsg ) ;
Lwm2mClientRpcRequest lwm2mClientRpcRequest = null ;
try {
Registration registration = clientContext . getClient ( new UUID ( sessionInfo . getSessionIdMSB ( ) , sessionInfo . getSessionIdLSB ( ) ) ) . getRegistration ( ) ;
lwm2mClientRpcRequest = this . getDeviceRpcRequest ( toDeviceRpcRequestMsg , sessionInfo , registration ) ;
if ( lwm2mClientRpcRequest . getErrorMsg ( ) ! = null ) {
lwm2mClientRpcRequest . setResponseCode ( BAD_REQUEST . name ( ) ) ;
this . onToDeviceRpcResponse ( lwm2mClientRpcRequest . getDeviceRpcResponseResultMsg ( ) , sessionInfo ) ;
} else {
lwM2mTransportRequest . sendAllRequest ( registration , lwm2mClientRpcRequest . getTargetIdVer ( ) , lwm2mClientRpcRequest . getTypeOper ( ) , lwm2mClientRpcRequest . getContentFormatName ( ) ,
lwm2mClientRpcRequest . getValue ( ) = = null ? lwm2mClientRpcRequest . getParams ( ) : lwm2mClientRpcRequest . getValue ( ) ,
this . config . getTimeout ( ) , lwm2mClientRpcRequest ) ;
}
} catch ( Exception e ) {
if ( lwm2mClientRpcRequest = = null ) {
lwm2mClientRpcRequest = new Lwm2mClientRpcRequest ( ) ;
}
lwm2mClientRpcRequest . setResponseCode ( BAD_REQUEST . name ( ) ) ;
if ( lwm2mClientRpcRequest . getErrorMsg ( ) = = null ) {
lwm2mClientRpcRequest . setErrorMsg ( e . getMessage ( ) ) ;
}
this . onToDeviceRpcResponse ( lwm2mClientRpcRequest . getDeviceRpcResponseResultMsg ( ) , sessionInfo ) ;
}
}
private Lwm2mClientRpcRequest getDeviceRpcRequest ( TransportProtos . ToDeviceRpcRequestMsg toDeviceRequest ,
SessionInfoProto sessionInfo , Registration registration ) throws IllegalArgumentException {
Lwm2mClientRpcRequest lwm2mClientRpcRequest = new Lwm2mClientRpcRequest ( ) ;
try {
lwm2mClientRpcRequest . setRequestId ( toDeviceRequest . getRequestId ( ) ) ;
lwm2mClientRpcRequest . setSessionInfo ( sessionInfo ) ;
lwm2mClientRpcRequest . setValidTypeOper ( toDeviceRequest . getMethodName ( ) ) ;
JsonObject rpcRequest = LwM2mTransportUtil . validateJson ( toDeviceRequest . getParams ( ) ) ;
if ( rpcRequest ! = null ) {
if ( rpcRequest . has ( lwm2mClientRpcRequest . keyNameKey ) ) {
String targetIdVer = this . getPresentPathIntoProfile ( sessionInfo ,
rpcRequest . get ( lwm2mClientRpcRequest . keyNameKey ) . getAsString ( ) ) ;
if ( targetIdVer ! = null ) {
lwm2mClientRpcRequest . setTargetIdVer ( targetIdVer ) ;
lwm2mClientRpcRequest . setInfoMsg ( String . format ( "Changed by: key - %s, pathIdVer - %s" ,
rpcRequest . get ( lwm2mClientRpcRequest . keyNameKey ) . getAsString ( ) , targetIdVer ) ) ;
}
}
if ( lwm2mClientRpcRequest . getTargetIdVer ( ) = = null ) {
lwm2mClientRpcRequest . setValidTargetIdVerKey ( rpcRequest , registration ) ;
}
if ( rpcRequest . has ( lwm2mClientRpcRequest . contentFormatNameKey ) ) {
lwm2mClientRpcRequest . setValidContentFormatName ( rpcRequest ) ;
}
if ( rpcRequest . has ( lwm2mClientRpcRequest . timeoutInMsKey ) & & rpcRequest . get ( lwm2mClientRpcRequest . timeoutInMsKey ) . getAsLong ( ) > 0 ) {
lwm2mClientRpcRequest . setTimeoutInMs ( rpcRequest . get ( lwm2mClientRpcRequest . timeoutInMsKey ) . getAsLong ( ) ) ;
}
if ( rpcRequest . has ( lwm2mClientRpcRequest . valueKey ) ) {
lwm2mClientRpcRequest . setValue ( rpcRequest . get ( lwm2mClientRpcRequest . valueKey ) . getAsString ( ) ) ;
}
if ( rpcRequest . has ( lwm2mClientRpcRequest . paramsKey ) & & rpcRequest . get ( lwm2mClientRpcRequest . paramsKey ) . isJsonObject ( ) ) {
ConcurrentHashMap < String , Object > params = new Gson ( ) . fromJson ( rpcRequest . get ( lwm2mClientRpcRequest . paramsKey )
. getAsJsonObject ( ) . toString ( ) , new TypeToken < ConcurrentHashMap < String , Object > > ( ) {
} . getType ( ) ) ;
if ( WRITE_UPDATE = = lwm2mClientRpcRequest . getTypeOper ( ) ) {
ConcurrentHashMap < String , Object > paramsResourceId = convertParamsToResourceId ( params , sessionInfo ) ;
if ( paramsResourceId . size ( ) > 0 ) {
lwm2mClientRpcRequest . setParams ( paramsResourceId ) ;
}
} else {
lwm2mClientRpcRequest . setParams ( params ) ;
}
} else if ( rpcRequest . has ( lwm2mClientRpcRequest . paramsKey ) & & rpcRequest . get ( lwm2mClientRpcRequest . paramsKey ) . isJsonArray ( ) ) {
new Gson ( ) . fromJson ( rpcRequest . get ( lwm2mClientRpcRequest . paramsKey )
. getAsJsonObject ( ) . toString ( ) , new TypeToken < ConcurrentHashMap < String , Object > > ( ) {
} . getType ( ) ) ;
// #1
this . checkRpcRequestTimeout ( ) ;
log . warn ( "4) toDeviceRpcRequestMsg: [{}], sessionUUID: [{}]" , toDeviceRpcRequestMsg , new UUID ( sessionInfo . getSessionIdMSB ( ) , sessionInfo . getSessionIdLSB ( ) ) ) ;
String bodyParams = StringUtils . trimToNull ( toDeviceRpcRequestMsg . getParams ( ) ) ! = null ? toDeviceRpcRequestMsg . getParams ( ) : "null" ;
LwM2mTypeOper lwM2mTypeOper = setValidTypeOper ( toDeviceRpcRequestMsg . getMethodName ( ) ) ;
UUID requestUUID = new UUID ( toDeviceRpcRequestMsg . getRequestIdMSB ( ) , toDeviceRpcRequestMsg . getRequestIdLSB ( ) ) ;
if ( ! this . rpcSubscriptions . containsKey ( requestUUID ) ) {
this . rpcSubscriptions . put ( requestUUID , toDeviceRpcRequestMsg . getExpirationTime ( ) ) ;
Lwm2mClientRpcRequest lwm2mClientRpcRequest = null ;
try {
Registration registration = clientContext . getClient ( sessionInfo ) . getRegistration ( ) ;
lwm2mClientRpcRequest = new Lwm2mClientRpcRequest ( lwM2mTypeOper , bodyParams , toDeviceRpcRequestMsg . getRequestId ( ) , sessionInfo , registration , this ) ;
if ( lwm2mClientRpcRequest . getErrorMsg ( ) ! = null ) {
lwm2mClientRpcRequest . setResponseCode ( BAD_REQUEST . name ( ) ) ;
this . onToDeviceRpcResponse ( lwm2mClientRpcRequest . getDeviceRpcResponseResultMsg ( ) , sessionInfo ) ;
} else {
lwM2mTransportRequest . sendAllRequest ( registration , lwm2mClientRpcRequest . getTargetIdVer ( ) , lwm2mClientRpcRequest . getTypeOper ( ) ,
null ,
lwm2mClientRpcRequest . getValue ( ) = = null ? lwm2mClientRpcRequest . getParams ( ) : lwm2mClientRpcRequest . getValue ( ) ,
this . config . getTimeout ( ) , lwm2mClientRpcRequest ) ;
}
lwm2mClientRpcRequest . setSessionInfo ( sessionInfo ) ;
if ( ! ( OBSERVE_READ_ALL = = lwm2mClientRpcRequest . getTypeOper ( )
| | DISCOVER_All = = lwm2mClientRpcRequest . getTypeOper ( )
| | OBSERVE_CANCEL = = lwm2mClientRpcRequest . getTypeOper ( ) )
& & lwm2mClientRpcRequest . getTargetIdVer ( ) = = null ) {
lwm2mClientRpcRequest . setErrorMsg ( lwm2mClientRpcRequest . targetIdVerKey + " and " +
lwm2mClientRpcRequest . keyNameKey + " is null or bad format" ) ;
} catch ( Exception e ) {
if ( lwm2mClientRpcRequest = = null ) {
lwm2mClientRpcRequest = new Lwm2mClientRpcRequest ( ) ;
}
/ * *
* EXECUTE & & WRITE_REPLACE - only for Resource or ResourceInstance
* /
else if ( ( EXECUTE = = lwm2mClientRpcRequest . getTypeOper ( )
| | WRITE_REPLACE = = lwm2mClientRpcRequest . getTypeOper ( ) )
& & lwm2mClientRpcRequest . getTargetIdVer ( ) ! = null
& & ! ( new LwM2mPath ( convertPathFromIdVerToObjectId ( lwm2mClientRpcRequest . getTargetIdVer ( ) ) ) . isResource ( )
| | new LwM2mPath ( convertPathFromIdVerToObjectId ( lwm2mClientRpcRequest . getTargetIdVer ( ) ) ) . isResourceInstance ( ) ) ) {
lwm2mClientRpcRequest . setErrorMsg ( "Invalid parameter " + lwm2mClientRpcRequest . targetIdVerKey
+ ". Only Resource or ResourceInstance can be this operation" ) ;
lwm2mClientRpcRequest . setResponseCode ( BAD_REQUEST . name ( ) ) ;
if ( lwm2mClientRpcRequest . getErrorMsg ( ) = = null ) {
lwm2mClientRpcRequest . setErrorMsg ( e . getMessage ( ) ) ;
}
} else {
lwm2mClientRpcRequest . setErrorMsg ( "Params of request is bad Json format." ) ;
this . onToDeviceRpcResponse ( lwm2mClientRpcRequest . getDeviceRpcResponseResultMsg ( ) , sessionInfo ) ;
}
} catch ( Exception e ) {
throw new IllegalArgumentException ( lwm2mClientRpcRequest . getErrorMsg ( ) ) ;
}
return lwm2mClientRpcRequest ;
}
private ConcurrentHashMap < String , Object > convertParamsToResourceId ( ConcurrentHashMap < String , Object > params ,
SessionInfoProto sessionInfo ) {
ConcurrentHashMap < String , Object > paramsIdVer = new ConcurrentHashMap < > ( ) ;
params . forEach ( ( k , v ) - > {
String targetIdVer = this . getPresentPathIntoProfile ( sessionInfo , k ) ;
if ( targetIdVer ! = null ) {
LwM2mPath targetId = new LwM2mPath ( convertPathFromIdVerToObjectId ( targetIdVer ) ) ;
if ( targetId . isResource ( ) ) {
paramsIdVer . put ( String . valueOf ( targetId . getResourceId ( ) ) , v ) ;
}
}
} ) ;
return paramsIdVer ;
private void checkRpcRequestTimeout ( ) {
Set < UUID > rpcSubscriptionsToRemove = rpcSubscriptions . entrySet ( ) . stream ( ) . filter ( kv - > System . currentTimeMillis ( ) > kv . getValue ( ) ) . map ( Map . Entry : : getKey ) . collect ( Collectors . toSet ( ) ) ;
rpcSubscriptionsToRemove . forEach ( rpcSubscriptions : : remove ) ;
}
public void sentRpcRequest ( Lwm2mClientRpcRequest rpcRequest , String requestCode , String msg , String typeMsg ) {
public void sentRpcResponse ( Lwm2mClientRpcRequest rpcRequest , String requestCode , String msg , String typeMsg ) {
rpcRequest . setResponseCode ( requestCode ) ;
if ( LOG_LW2M_ERROR . equals ( typeMsg ) ) {
rpcRequest . setInfoMsg ( null ) ;
@ -578,6 +506,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
@Override
public void onToDeviceRpcResponse ( TransportProtos . ToDeviceRpcResponseMsg toDeviceResponse , SessionInfoProto sessionInfo ) {
log . warn ( "5) onToDeviceRpcResponse: [{}], sessionUUID: [{}]" , toDeviceResponse , new UUID ( sessionInfo . getSessionIdMSB ( ) , sessionInfo . getSessionIdLSB ( ) ) ) ;
transportService . process ( sessionInfo , toDeviceResponse , null ) ;
}
@ -722,10 +651,10 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
* set setClient_fw_info . . . = value
* * /
if ( lwM2MClient . getFwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getFwUpdate ( ) . initReadValue ( this , lwM2mTransportRequest , path ) ;
lwM2MClient . getFwUpdate ( ) . initReadValue ( this , this . lwM2mTransportRequest , path ) ;
}
if ( lwM2MClient . getSwUpdate ( ) . isInfoFwSwUpdate ( ) ) {
lwM2MClient . getSwUpdate ( ) . initReadValue ( this , lwM2mTransportRequest , path ) ;
lwM2MClient . getSwUpdate ( ) . initReadValue ( this , this . lwM2mTransportRequest , path ) ;
}
/ * *
@ -742,7 +671,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
& & ( convertPathFromObjectIdToIdVer ( FW_RESULT_ID , registration ) . equals ( path ) ) ) {
if ( DOWNLOADED . name ( ) . equals ( lwM2MClient . getFwUpdate ( ) . getStateUpdate ( ) )
& & lwM2MClient . getFwUpdate ( ) . conditionalFwExecuteStart ( ) ) {
lwM2MClient . getFwUpdate ( ) . executeFwSwWare ( this , lwM2mTransportRequest ) ;
lwM2MClient . getFwUpdate ( ) . executeFwSwWare ( this , this . lwM2mTransportRequest ) ;
} else if ( UPDATING . name ( ) . equals ( lwM2MClient . getFwUpdate ( ) . getStateUpdate ( ) )
& & lwM2MClient . getFwUpdate ( ) . conditionalFwExecuteAfterSuccess ( ) ) {
lwM2MClient . getFwUpdate ( ) . finishFwSwUpdate ( this , true ) ;
@ -767,7 +696,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
& & ( convertPathFromObjectIdToIdVer ( SW_RESULT_ID , registration ) . equals ( path ) ) ) {
if ( DOWNLOADED . name ( ) . equals ( lwM2MClient . getSwUpdate ( ) . getStateUpdate ( ) )
& & lwM2MClient . getSwUpdate ( ) . conditionalSwUpdateExecute ( ) ) {
lwM2MClient . getSwUpdate ( ) . executeFwSwWare ( this , lwM2mTransportRequest ) ;
lwM2MClient . getSwUpdate ( ) . executeFwSwWare ( this , this . lwM2mTransportRequest ) ;
} else if ( UPDATING . name ( ) . equals ( lwM2MClient . getSwUpdate ( ) . getStateUpdate ( ) )
& & lwM2MClient . getSwUpdate ( ) . conditionalSwExecuteAfterSuccess ( ) ) {
lwM2MClient . getSwUpdate ( ) . finishFwSwUpdate ( this , true ) ;
@ -958,10 +887,16 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
* /
private Object getResourceValueFormatKv ( LwM2mClient lwM2MClient , String pathIdVer ) {
LwM2mResource resourceValue = this . getResourceValueFromLwM2MClient ( lwM2MClient , pathIdVer ) ;
ResourceModel . Type currentType = resourceValue . getType ( ) ;
ResourceModel . Type expectedType = this . helper . getResourceModelTypeEqualsKvProtoValueType ( currentType , pathIdVer ) ;
return this . converter . convertValue ( resourceValue . getValue ( ) , currentType , expectedType ,
new LwM2mPath ( convertPathFromIdVerToObjectId ( pathIdVer ) ) ) ;
if ( resourceValue ! = null ) {
ResourceModel . Type currentType = resourceValue . getType ( ) ;
ResourceModel . Type expectedType = this . helper . getResourceModelTypeEqualsKvProtoValueType ( currentType , pathIdVer ) ;
return this . converter . convertValue ( resourceValue . getValue ( ) , currentType , expectedType ,
new LwM2mPath ( convertPathFromIdVerToObjectId ( pathIdVer ) ) ) ;
}
else {
return null ;
}
}
/ * *
@ -970,11 +905,14 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
* @return - return value of Resource by idPath
* /
private LwM2mResource getResourceValueFromLwM2MClient ( LwM2mClient lwM2MClient , String path ) {
LwM2mResource resourceValue = null ;
if ( new LwM2mPath ( convertPathFromIdVerToObjectId ( path ) ) . isResource ( ) ) {
resourceValue = lwM2MClient . getResources ( ) . get ( path ) . getLwM2mResource ( ) ;
LwM2mResource lwm2mResourceValue = null ;
ResourceValue resourceValue = lwM2MClient . getResources ( ) . get ( path ) ;
if ( resourceValue ! = null ) {
if ( new LwM2mPath ( convertPathFromIdVerToObjectId ( path ) ) . isResource ( ) ) {
lwm2mResourceValue = lwM2MClient . getResources ( ) . get ( path ) . getLwM2mResource ( ) ;
}
}
return resourceValue ;
return lwm2mR esourceValue;
}
/ * *
@ -1281,7 +1219,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
* @param name -
* @return -
* /
private String getPresentPathIntoProfile ( TransportProtos . SessionInfoProto sessionInfo , String name ) {
public String getPresentPathIntoProfile ( TransportProtos . SessionInfoProto sessionInfo , String name ) {
LwM2mClientProfile profile = clientContext . getProfile ( new UUID ( sessionInfo . getDeviceProfileIdMSB ( ) , sessionInfo . getDeviceProfileIdLSB ( ) ) ) ;
LwM2mClient lwM2mClient = clientContext . getClient ( sessionInfo ) ;
return profile . getPostKeyNameProfile ( ) . getAsJsonObject ( ) . entrySet ( ) . stream ( )
@ -1301,10 +1239,9 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
public void onGetAttributesResponse ( TransportProtos . GetAttributeResponseMsg attributesResponse , TransportProtos . SessionInfoProto sessionInfo ) {
try {
List < TransportProtos . TsKvProto > tsKvProtos = attributesResponse . getSharedAttributeListList ( ) ;
this . updateAttributeFromThingsboard ( tsKvProtos , sessionInfo ) ;
} catch ( Exception e ) {
log . error ( String . valueOf ( e ) ) ;
log . error ( "" , e ) ;
}
}
@ -1320,22 +1257,28 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
* /
public void updateAttributeFromThingsboard ( List < TransportProtos . TsKvProto > tsKvProtos , TransportProtos . SessionInfoProto sessionInfo ) {
LwM2mClient lwM2MClient = clientContext . getClient ( sessionInfo ) ;
tsKvProtos . forEach ( tsKvProto - > {
String pathIdVer = this . getPresentPathIntoProfile ( sessionInfo , tsKvProto . getKv ( ) . getKey ( ) ) ;
if ( pathIdVer ! = null ) {
// #1.1
if ( lwM2MClient . getDelayedRequests ( ) . containsKey ( pathIdVer ) & & tsKvProto . getTs ( ) > lwM2MClient . getDelayedRequests ( ) . get ( pathIdVer ) . getTs ( ) ) {
lwM2MClient . getDelayedRequests ( ) . put ( pathIdVer , tsKvProto ) ;
} else if ( ! lwM2MClient . getDelayedRequests ( ) . containsKey ( pathIdVer ) ) {
lwM2MClient . getDelayedRequests ( ) . put ( pathIdVer , tsKvProto ) ;
if ( lwM2MClient ! = null ) {
log . warn ( "1) UpdateAttributeFromThingsboard, tsKvProtos [{}]" , tsKvProtos ) ;
tsKvProtos . forEach ( tsKvProto - > {
String pathIdVer = this . getPresentPathIntoProfile ( sessionInfo , tsKvProto . getKv ( ) . getKey ( ) ) ;
if ( pathIdVer ! = null ) {
// #1.1
if ( lwM2MClient . getDelayedRequests ( ) . containsKey ( pathIdVer ) & & tsKvProto . getTs ( ) > lwM2MClient . getDelayedRequests ( ) . get ( pathIdVer ) . getTs ( ) ) {
lwM2MClient . getDelayedRequests ( ) . put ( pathIdVer , tsKvProto ) ;
} else if ( ! lwM2MClient . getDelayedRequests ( ) . containsKey ( pathIdVer ) ) {
lwM2MClient . getDelayedRequests ( ) . put ( pathIdVer , tsKvProto ) ;
}
}
}
} ) ;
// #2.1
lwM2MClient . getDelayedRequests ( ) . forEach ( ( pathIdVer , tsKvProto ) - > {
this . updateResourcesValueToClient ( lwM2MClient , this . getResourceValueFormatKv ( lwM2MClient , pathIdVer ) ,
getValueFromKvProto ( tsKvProto . getKv ( ) ) , pathIdVer ) ;
} ) ;
} ) ;
// #2.1
lwM2MClient . getDelayedRequests ( ) . forEach ( ( pathIdVer , tsKvProto ) - > {
this . updateResourcesValueToClient ( lwM2MClient , this . getResourceValueFormatKv ( lwM2MClient , pathIdVer ) ,
getValueFromKvProto ( tsKvProto . getKv ( ) ) , pathIdVer ) ;
} ) ;
}
else {
log . error ( "UpdateAttributeFromThingsboard, lwM2MClient is null" ) ;
}
}
/ * *
@ -1378,6 +1321,7 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
private void reportActivityAndRegister ( SessionInfoProto sessionInfo ) {
if ( sessionInfo ! = null & & transportService . reportActivity ( sessionInfo ) = = null ) {
transportService . registerAsyncSession ( sessionInfo , new LwM2mSessionMsgListener ( this , sessionInfo ) ) ;
this . reportActivitySubscription ( sessionInfo ) ;
}
}
@ -1413,23 +1357,30 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
}
}
public void getInfoFirmwareUpdate ( LwM2mClient lwM2MClient ) {
public void getInfoFirmwareUpdate ( LwM2mClient lwM2MClient , Lwm2mClientRpcRequest rpcRequest ) {
if ( lwM2MClient . getRegistration ( ) . getSupportedVersion ( FW_ID ) ! = null ) {
SessionInfoProto sessionInfo = this . getSessionInfoOrCloseSession ( lwM2MClient ) ;
if ( sessionInfo ! = null ) {
DefaultLwM2MTransportMsgHandler serviceImpl = this ;
transportService . process ( sessionInfo , createFirmwar eRequestMsg ( sessionInfo , Firmwar eType. FIRMWARE . name ( ) ) ,
DefaultLwM2MTransportMsgHandler handler = this ;
this . transportService . process ( sessionInfo , createOtaPackag eRequestMsg ( sessionInfo , OtaPackag eType. FIRMWARE . name ( ) ) ,
new TransportServiceCallback < > ( ) {
@Override
public void onSuccess ( TransportProtos . GetFirmwar eResponseMsg response ) {
public void onSuccess ( TransportProtos . GetOtaPackag eResponseMsg response ) {
if ( TransportProtos . ResponseStatus . SUCCESS . equals ( response . getResponseStatus ( ) )
& & response . getType ( ) . equals ( FirmwareType . FIRMWARE . name ( ) ) ) {
& & response . getType ( ) . equals ( OtaPackageType . FIRMWARE . name ( ) ) ) {
log . warn ( "7) firmware start with ver: [{}]" , response . getVersion ( ) ) ;
lwM2MClient . getFwUpdate ( ) . setRpcRequest ( rpcRequest ) ;
lwM2MClient . getFwUpdate ( ) . setCurrentVersion ( response . getVersion ( ) ) ;
lwM2MClient . getFwUpdate ( ) . setCurrentTitle ( response . getTitle ( ) ) ;
lwM2MClient . getFwUpdate ( ) . setCurrentId ( new FirmwareId ( new UUID ( response . getFirmwareIdMSB ( ) , response . getFirmwareIdLSB ( ) ) ) . getId ( ) ) ;
lwM2MClient . getFwUpdate ( ) . sendReadObserveInfo ( lwM2mTransportRequest ) ;
lwM2MClient . getFwUpdate ( ) . setCurrentId ( new OtaPackageId ( new UUID ( response . getOtaPackageIdMSB ( ) , response . getOtaPackageIdLSB ( ) ) ) . getId ( ) ) ;
if ( rpcRequest = = null ) {
lwM2MClient . getFwUpdate ( ) . sendReadObserveInfo ( lwM2mTransportRequest ) ;
}
else {
lwM2MClient . getFwUpdate ( ) . writeFwSwWare ( handler , lwM2mTransportRequest ) ;
}
} else {
log . trace ( "Firmware [{}] [{}]" , lwM2MClient . getDeviceName ( ) , response . getResponseStatus ( ) . toString ( ) ) ;
log . trace ( "OtaPackag e [{}] [{}]" , lwM2MClient . getDeviceName ( ) , response . getResponseStatus ( ) . toString ( ) ) ;
}
}
@ -1442,21 +1393,28 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
}
}
public void getInfoSoftwareUpdate ( LwM2mClient lwM2MClient ) {
public void getInfoSoftwareUpdate ( LwM2mClient lwM2MClient , Lwm2mClientRpcRequest rpcRequest ) {
if ( lwM2MClient . getRegistration ( ) . getSupportedVersion ( SW_ID ) ! = null ) {
SessionInfoProto sessionInfo = this . getSessionInfoOrCloseSession ( lwM2MClient ) ;
if ( sessionInfo ! = null ) {
DefaultLwM2MTransportMsgHandler serviceImpl = this ;
transportService . process ( sessionInfo , createFirmwar eRequestMsg ( sessionInfo , Firmwar eType. SOFTWARE . name ( ) ) ,
DefaultLwM2MTransportMsgHandler handler = this ;
transportService . process ( sessionInfo , createOtaPackag eRequestMsg ( sessionInfo , OtaPackag eType. SOFTWARE . name ( ) ) ,
new TransportServiceCallback < > ( ) {
@Override
public void onSuccess ( TransportProtos . GetFirmwar eResponseMsg response ) {
public void onSuccess ( TransportProtos . GetOtaPackag eResponseMsg response ) {
if ( TransportProtos . ResponseStatus . SUCCESS . equals ( response . getResponseStatus ( ) )
& & response . getType ( ) . equals ( FirmwareType . SOFTWARE . name ( ) ) ) {
& & response . getType ( ) . equals ( OtaPackageType . SOFTWARE . name ( ) ) ) {
lwM2MClient . getSwUpdate ( ) . setRpcRequest ( rpcRequest ) ;
lwM2MClient . getSwUpdate ( ) . setCurrentVersion ( response . getVersion ( ) ) ;
lwM2MClient . getSwUpdate ( ) . setCurrentTitle ( response . getTitle ( ) ) ;
lwM2MClient . getSwUpdate ( ) . setCurrentId ( new Firmwar eId( new UUID ( response . getFirmwar eIdMSB ( ) , response . getFirmwar eIdLSB ( ) ) ) . getId ( ) ) ;
lwM2MClient . getSwUpdate ( ) . setCurrentId ( new OtaPackag eId( new UUID ( response . getOtaPackag eIdMSB ( ) , response . getOtaPackag eIdLSB ( ) ) ) . getId ( ) ) ;
lwM2MClient . getSwUpdate ( ) . sendReadObserveInfo ( lwM2mTransportRequest ) ;
if ( rpcRequest = = null ) {
lwM2MClient . getSwUpdate ( ) . sendReadObserveInfo ( lwM2mTransportRequest ) ;
}
else {
lwM2MClient . getSwUpdate ( ) . writeFwSwWare ( handler , lwM2mTransportRequest ) ;
}
} else {
log . trace ( "Software [{}] [{}]" , lwM2MClient . getDeviceName ( ) , response . getResponseStatus ( ) . toString ( ) ) ;
}
@ -1471,8 +1429,8 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
}
}
private TransportProtos . GetFirmwareRequestMsg createFirmwar eRequestMsg ( SessionInfoProto sessionInfo , String nameFwSW ) {
return TransportProtos . GetFirmwar eRequestMsg . newBuilder ( )
private TransportProtos . GetOtaPackageRequestMsg createOtaPackag eRequestMsg ( SessionInfoProto sessionInfo , String nameFwSW ) {
return TransportProtos . GetOtaPackag eRequestMsg . newBuilder ( )
. setDeviceIdMSB ( sessionInfo . getDeviceIdMSB ( ) )
. setDeviceIdLSB ( sessionInfo . getDeviceIdLSB ( ) )
. setTenantIdMSB ( sessionInfo . getTenantIdMSB ( ) )
@ -1510,4 +1468,11 @@ public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler
return this . config ;
}
private void reportActivitySubscription ( TransportProtos . SessionInfoProto sessionInfo ) {
transportService . process ( sessionInfo , TransportProtos . SubscriptionInfoProto . newBuilder ( )
. setAttributeSubscription ( true )
. setRpcSubscription ( true )
. setLastActivityTime ( System . currentTimeMillis ( ) )
. build ( ) , TransportServiceCallback . EMPTY ) ;
}
}