@ -5,7 +5,7 @@
* you may not use this file except in compliance with the License .
* You may obtain a copy of the License at
*
* http : //www.apache.org/licenses/LICENSE-2.0
* http : //www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing , software
* distributed under the License is distributed on an "AS IS" BASIS ,
@ -48,15 +48,11 @@ import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTrans
import org.thingsboard.server.common.data.ota.OtaPackageUtil ;
import org.thingsboard.server.common.transport.TransportService ;
import org.thingsboard.server.common.transport.TransportServiceCallback ;
import org.thingsboard.server.common.transport.service.DefaultTransportService ;
import org.thingsboard.server.gen.transport.TransportProtos ;
import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent ;
import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto ;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent ;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig ;
import org.thingsboard.server.transport.lwm2m.server.LwM2mOtaConvert ;
import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest ;
import org.thingsboard.server.transport.lwm2m.server.LwM2mSessionMsgListener ;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext ;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper ;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil ;
@ -84,6 +80,7 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttrib
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService ;
import org.thingsboard.server.transport.lwm2m.server.ota.LwM2MOtaUpdateService ;
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcRequestHandler ;
import org.thingsboard.server.transport.lwm2m.server.session.LwM2MSessionManager ;
import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore ;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl ;
@ -129,6 +126,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
private final TransportService transportService ;
private final LwM2mTransportContext context ;
private final LwM2MAttributesService attributesService ;
private final LwM2MSessionManager sessionManager ;
private final LwM2MOtaUpdateService otaService ;
private final LwM2MTransportServerConfig config ;
private final LwM2MTelemetryLogService logService ;
@ -143,12 +141,14 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
LwM2mTransportServerHelper helper ,
LwM2mClientContext clientContext ,
LwM2MTelemetryLogService logService ,
LwM2MSessionManager sessionManager ,
@Lazy LwM2MOtaUpdateService otaService ,
@Lazy LwM2MAttributesService attributesService ,
@Lazy LwM2MRpcRequestHandler rpcHandler ,
@Lazy LwM2mDownlinkMsgHandler defaultLwM2MDownlinkMsgHandler ,
LwM2mTransportContext context , TbLwM2MDtlsSessionStore sessionStore ) {
this . transportService = transportService ;
this . sessionManager = sessionManager ;
this . attributesService = attributesService ;
this . otaService = otaService ;
this . config = config ;
@ -205,18 +205,10 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
Optional < SessionInfoProto > oldSessionInfo = this . clientContext . register ( lwM2MClient , registration ) ;
if ( oldSessionInfo . isPresent ( ) ) {
log . info ( "[{}] Closing old session: {}" , registration . getEndpoint ( ) , new UUID ( oldSessionInfo . get ( ) . getSessionIdMSB ( ) , oldSessionInfo . get ( ) . getSessionIdLSB ( ) ) ) ;
clo seS ession( oldSessionInfo . get ( ) ) ;
sessionManager . deregister ( oldSessionInfo . get ( ) ) ;
}
logService . log ( lwM2MClient , LOG_LWM2M_INFO + ": Client registered with registration id: " + registration . getId ( ) ) ;
SessionInfoProto sessionInfo = lwM2MClient . getSession ( ) ;
transportService . registerAsyncSession ( sessionInfo , new LwM2mSessionMsgListener ( this , attributesService , rpcHandler , sessionInfo , transportService ) ) ;
TransportProtos . TransportToDeviceActorMsg msg = TransportProtos . TransportToDeviceActorMsg . newBuilder ( )
. setSessionInfo ( sessionInfo )
. setSessionEvent ( DefaultTransportService . getSessionEventMsg ( SessionEvent . OPEN ) )
. setSubscribeToAttributes ( TransportProtos . SubscribeToAttributeUpdatesMsg . newBuilder ( ) . setSessionType ( TransportProtos . SessionType . ASYNC ) . build ( ) )
. setSubscribeToRPC ( TransportProtos . SubscribeToRPCMsg . newBuilder ( ) . setSessionType ( TransportProtos . SessionType . ASYNC ) . build ( ) )
. build ( ) ;
transportService . process ( msg , null ) ;
sessionManager . register ( lwM2MClient . getSession ( ) ) ;
this . initClientTelemetry ( lwM2MClient ) ;
this . initAttributes ( lwM2MClient ) ;
otaService . init ( lwM2MClient ) ;
@ -273,7 +265,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
clientContext . unregister ( client , registration ) ;
SessionInfoProto sessionInfo = client . getSession ( ) ;
if ( sessionInfo ! = null ) {
clo seS ession( sessionInfo ) ;
sessionManager . deregister ( sessionInfo ) ;
sessionStore . remove ( registration . getEndpoint ( ) ) ;
log . info ( "Client close session: [{}] unReg [{}] name [{}] profile " , registration . getId ( ) , registration . getEndpoint ( ) , sessionInfo . getDeviceType ( ) ) ;
} else {
@ -288,11 +280,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
} ) ;
}
public void closeSession ( SessionInfoProto sessionInfo ) {
transportService . process ( sessionInfo , DefaultTransportService . getSessionEventMsg ( SessionEvent . CLOSED ) , null ) ;
transportService . deregisterSession ( sessionInfo ) ;
}
@Override
public void onSleepingDev ( Registration registration ) {
log . info ( "[{}] [{}] Received endpoint Sleeping version event" , registration . getId ( ) , registration . getEndpoint ( ) ) ;
@ -300,19 +287,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
//TODO: associate endpointId with device information.
}
// /**
// * Cancel observation for All objects for this registration
// */
// @Override
// public void setCancelObservationsAll(Registration registration) {
// if (registration != null) {
// LwM2mClient client = clientContext.getClientByEndpoint(registration.getEndpoint());
// if (client != null && client.getRegistration() != null && client.getRegistration().getId().equals(registration.getId())) {
// defaultLwM2MDownlinkMsgHandler.sendCancelAllRequest(client, TbLwM2MCancelAllRequest.builder().build(), new TbLwM2MCancelAllObserveRequestCallback(this, client));
// }
// }
// }
/ * *
* Sending observe value to thingsboard from ObservationListener . onResponse : object , instance , SingleResource or MultipleResource
*
@ -337,6 +311,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
this . updateResourcesValue ( lwM2MClient , lwM2mResource , path ) ;
}
}
clientContext . update ( lwM2MClient ) ;
}
}
@ -375,16 +350,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
clientContext . getLwM2mClients ( ) . forEach ( e - > e . deleteResources ( pathIdVer , this . config . getModelProvider ( ) ) ) ;
}
/ * *
* Deregister session in transport
*
* @param sessionInfo - lwm2m client
* /
@Override
public void doDisconnect ( SessionInfoProto sessionInfo ) {
closeSession ( sessionInfo ) ;
}
/ * *
* Those methods are called by the protocol stage thread pool , this means that execution MUST be done in a short delay ,
* * if you need to do long time processing use a dedicated thread pool .
@ -479,14 +444,6 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
attributesMap . forEach ( ( targetId , params ) - > sendWriteAttributesRequest ( lwM2MClient , targetId , params ) ) ;
}
private void sendDiscoverRequests ( LwM2mClient lwM2MClient , Lwm2mDeviceProfileTransportConfiguration profile , Set < String > supportedObjects ) {
Set < String > targetIds = profile . getObserveAttr ( ) . getAttributeLwm2m ( ) . keySet ( ) ;
targetIds = targetIds . stream ( ) . filter ( target - > isSupportedTargetId ( supportedObjects , target ) ) . collect ( Collectors . toSet ( ) ) ;
// TODO: why do we need to put observe into pending read requests?
// lwM2MClient.getPendingReadRequests().addAll(targetIds);
targetIds . forEach ( targetId - > sendDiscoverRequest ( lwM2MClient , targetId ) ) ;
}
private void sendDiscoverRequest ( LwM2mClient lwM2MClient , String targetId ) {
TbLwM2MDiscoverRequest request = TbLwM2MDiscoverRequest . builder ( ) . versionedId ( targetId ) . timeout ( this . config . getTimeout ( ) ) . build ( ) ;
defaultLwM2MDownlinkMsgHandler . sendDiscoverRequest ( lwM2MClient , request , new TbLwM2MDiscoverCallback ( logService , lwM2MClient , targetId ) ) ;
@ -637,7 +594,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
List < TransportProtos . KeyValueProto > resultAttributes = new ArrayList < > ( ) ;
profile . getObserveAttr ( ) . getAttribute ( ) . forEach ( pathIdVer - > {
if ( path . contains ( pathIdVer ) ) {
TransportProtos . KeyValueProto kvAttr = this . getKvToThingsb oard ( pathIdVer , registration ) ;
TransportProtos . KeyValueProto kvAttr = this . getKvToThingsB oard ( pathIdVer , registration ) ;
if ( kvAttr ! = null ) {
resultAttributes . add ( kvAttr ) ;
}
@ -646,7 +603,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
List < TransportProtos . KeyValueProto > resultTelemetries = new ArrayList < > ( ) ;
profile . getObserveAttr ( ) . getTelemetry ( ) . forEach ( pathIdVer - > {
if ( path . contains ( pathIdVer ) ) {
TransportProtos . KeyValueProto kvAttr = this . getKvToThingsb oard ( pathIdVer , registration ) ;
TransportProtos . KeyValueProto kvAttr = this . getKvToThingsB oard ( pathIdVer , registration ) ;
if ( kvAttr ! = null ) {
resultTelemetries . add ( kvAttr ) ;
}
@ -663,7 +620,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
return null ;
}
private TransportProtos . KeyValueProto getKvToThingsb oard ( String pathIdVer , Registration registration ) {
private TransportProtos . KeyValueProto getKvToThingsB oard ( String pathIdVer , Registration registration ) {
LwM2mClient lwM2MClient = this . clientContext . getClientByEndpoint ( registration . getEndpoint ( ) ) ;
Map < String , String > names = clientContext . getProfile ( lwM2MClient . getProfileId ( ) ) . getObserveAttr ( ) . getKeyName ( ) ;
if ( names ! = null & & names . containsKey ( pathIdVer ) ) {
@ -710,10 +667,12 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
public void onWriteResponseOk ( LwM2mClient client , String path , WriteRequest request ) {
if ( request . getNode ( ) instanceof LwM2mResource ) {
this . updateResourcesValue ( client , ( ( LwM2mResource ) request . getNode ( ) ) , path ) ;
clientContext . update ( client ) ;
} else if ( request . getNode ( ) instanceof LwM2mObjectInstance ) {
( ( LwM2mObjectInstance ) request . getNode ( ) ) . getResources ( ) . forEach ( ( resId , resource ) - > {
this . updateResourcesValue ( client , resource , path + "/" + resId ) ;
} ) ;
clientContext . update ( client ) ;
}
}
@ -788,7 +747,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
if ( ! newLwM2mSettings . getFwUpdateStrategy ( ) . equals ( oldLwM2mSettings . getFwUpdateStrategy ( ) )
| | ( StringUtils . isNotEmpty ( newLwM2mSettings . getFwUpdateResource ( ) ) & &
! newLwM2mSettings . getFwUpdateResource ( ) . equals ( oldLwM2mSettings . getFwUpdateResource ( ) ) ) ) {
clients . forEach ( lwM2MClient - > otaService . onCurrent FirmwareStrategyUpdate ( lwM2MClient , newLwM2mSettings ) ) ;
clients . forEach ( lwM2MClient - > otaService . onFirmwareStrategyUpdate ( lwM2MClient , newLwM2mSettings ) ) ;
}
if ( ! newLwM2mSettings . getSwUpdateStrategy ( ) . equals ( oldLwM2mSettings . getSwUpdateStrategy ( ) )
@ -893,7 +852,7 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl
* /
private void reportActivityAndRegister ( SessionInfoProto sessionInfo ) {
if ( sessionInfo ! = null & & transportService . reportActivity ( sessionInfo ) = = null ) {
transportService . registerAsyncSession ( sessionInfo , new LwM2mSessionMsgListener ( this , attributesService , rpcHandler , sessionInfo , transportService ) ) ;
sessionManager . register ( sessionInfo ) ;
this . reportActivitySubscription ( sessionInfo ) ;
}
}