Browse Source

Refactoring of LwM2M transport

pull/4536/head
Andrii Shvaika 5 years ago
parent
commit
8b3e34f0ef
  1. 21
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapServerConfiguration.java
  2. 15
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java
  3. 17
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java
  4. 105
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2MTransportMsgHandler.java
  5. 77
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java
  6. 22
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportService.java
  7. 4
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java
  8. 4
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java
  9. 32
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportContext.java
  10. 6
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportMsgHandler.java
  11. 37
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportRequest.java
  12. 83
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerHelper.java
  13. 31
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerInitializer.java
  14. 16
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mVersionedModelProvider.java
  15. 4
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java
  16. 4
      common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportContext.java

21
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/LwM2MTransportBootstrapServerConfiguration.java

@ -30,7 +30,8 @@ import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2MBootstrapSec
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2MInMemoryBootstrapConfigStore;
import org.thingsboard.server.transport.lwm2m.bootstrap.secure.LwM2mDefaultBootstrapSessionManager;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportBootstrapConfig;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContextServer;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper;
import java.math.BigInteger;
import java.security.AlgorithmParameters;
@ -68,10 +69,10 @@ public class LwM2MTransportBootstrapServerConfiguration {
private boolean pskMode = false;
@Autowired
private LwM2MTransportContextBootstrap contextBs;
private LwM2MTransportServerConfig serverConfig;
@Autowired
private LwM2mTransportContextServer contextS;
private LwM2MTransportContextBootstrap contextBs;
@Autowired
private LwM2MBootstrapSecurityStore lwM2MBootstrapSecurityStore;
@ -108,8 +109,8 @@ public class LwM2MTransportBootstrapServerConfiguration {
/** Create and Set DTLS Config */
DtlsConnectorConfig.Builder dtlsConfig = new DtlsConnectorConfig.Builder();
dtlsConfig.setRecommendedSupportedGroupsOnly(this.contextS.getLwM2MTransportServerConfig().isRecommendedSupportedGroups());
dtlsConfig.setRecommendedCipherSuitesOnly(this.contextS.getLwM2MTransportServerConfig().isRecommendedCiphers());
dtlsConfig.setRecommendedSupportedGroupsOnly(serverConfig.isRecommendedSupportedGroups());
dtlsConfig.setRecommendedCipherSuitesOnly(serverConfig.isRecommendedCiphers());
if (this.pskMode) {
dtlsConfig.setSupportedCipherSuites(
TLS_PSK_WITH_AES_128_CCM_8,
@ -134,10 +135,10 @@ public class LwM2MTransportBootstrapServerConfiguration {
private void setServerWithCredentials(LeshanBootstrapServerBuilder builder) {
try {
if (this.contextS.getLwM2MTransportServerConfig().getKeyStoreValue() != null) {
KeyStore keyStoreServer = this.contextS.getLwM2MTransportServerConfig().getKeyStoreValue();
if (serverConfig.getKeyStoreValue() != null) {
KeyStore keyStoreServer = serverConfig.getKeyStoreValue();
if (this.setBuilderX509(builder)) {
X509Certificate rootCAX509Cert = (X509Certificate) keyStoreServer.getCertificate(this.contextS.getLwM2MTransportServerConfig().getRootCertificateAlias());
X509Certificate rootCAX509Cert = (X509Certificate) keyStoreServer.getCertificate(serverConfig.getRootCertificateAlias());
if (rootCAX509Cert != null) {
X509Certificate[] trustedCertificates = new X509Certificate[1];
trustedCertificates[0] = rootCAX509Cert;
@ -168,8 +169,8 @@ public class LwM2MTransportBootstrapServerConfiguration {
* For idea => KeyStorePathResource == common/transport/lwm2m/src/main/resources/credentials: in LwM2MTransportContextServer: credentials/serverKeyStore.jks
*/
try {
X509Certificate serverCertificate = (X509Certificate) this.contextS.getLwM2MTransportServerConfig().getKeyStoreValue().getCertificate(this.contextBs.getCtxBootStrap().getCertificateAlias());
PrivateKey privateKey = (PrivateKey) this.contextS.getLwM2MTransportServerConfig().getKeyStoreValue().getKey(this.contextBs.getCtxBootStrap().getCertificateAlias(), this.contextS.getLwM2MTransportServerConfig().getKeyStorePassword() == null ? null : this.contextS.getLwM2MTransportServerConfig().getKeyStorePassword().toCharArray());
X509Certificate serverCertificate = (X509Certificate) serverConfig.getKeyStoreValue().getCertificate(this.contextBs.getCtxBootStrap().getCertificateAlias());
PrivateKey privateKey = (PrivateKey) serverConfig.getKeyStoreValue().getKey(this.contextBs.getCtxBootStrap().getCertificateAlias(), serverConfig.getKeyStorePassword() == null ? null : serverConfig.getKeyStorePassword().toCharArray());
PublicKey publicKey = serverCertificate.getPublicKey();
if (privateKey != null && privateKey.getEncoded().length > 0 && publicKey != null && publicKey.getEncoded().length > 0) {
builder.setPublicKey(serverCertificate.getPublicKey());

15
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java

@ -34,7 +34,8 @@ import org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode;
import org.thingsboard.server.transport.lwm2m.secure.LwM2mCredentialsSecurityInfoValidator;
import org.thingsboard.server.transport.lwm2m.secure.ReadResultSecurityStore;
import org.thingsboard.server.transport.lwm2m.server.LwM2mSessionMsgListener;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContextServer;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandlerUtil;
import java.io.IOException;
@ -59,12 +60,14 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
private final LwM2mCredentialsSecurityInfoValidator lwM2MCredentialsSecurityInfoValidator;
private final LwM2mTransportContextServer context;
private final LwM2mTransportContext context;
private final LwM2mTransportServerHelper helper;
public LwM2MBootstrapSecurityStore(EditableBootstrapConfigStore bootstrapConfigStore, LwM2mCredentialsSecurityInfoValidator lwM2MCredentialsSecurityInfoValidator, LwM2mTransportContextServer context) {
public LwM2MBootstrapSecurityStore(EditableBootstrapConfigStore bootstrapConfigStore, LwM2mCredentialsSecurityInfoValidator lwM2MCredentialsSecurityInfoValidator, LwM2mTransportContext context, LwM2mTransportServerHelper helper) {
this.bootstrapConfigStore = bootstrapConfigStore;
this.lwM2MCredentialsSecurityInfoValidator = lwM2MCredentialsSecurityInfoValidator;
this.context = context;
this.helper = helper;
}
@Override
@ -158,19 +161,19 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
LwM2MServerBootstrap profileServerBootstrap = mapper.readValue(bootstrapObject.get(BOOTSTRAP_SERVER).toString(), LwM2MServerBootstrap.class);
LwM2MServerBootstrap profileLwm2mServer = mapper.readValue(bootstrapObject.get(LWM2M_SERVER).toString(), LwM2MServerBootstrap.class);
UUID sessionUUiD = UUID.randomUUID();
TransportProtos.SessionInfoProto sessionInfo = context.getValidateSessionInfo(store.getMsg(), sessionUUiD.getMostSignificantBits(), sessionUUiD.getLeastSignificantBits());
TransportProtos.SessionInfoProto sessionInfo = helper.getValidateSessionInfo(store.getMsg(), sessionUUiD.getMostSignificantBits(), sessionUUiD.getLeastSignificantBits());
context.getTransportService().registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(null, sessionInfo));
if (this.getValidatedSecurityMode(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap, lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer)) {
lwM2MBootstrapConfig.bootstrapServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap);
lwM2MBootstrapConfig.lwm2mServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer);
String logMsg = String.format("%s: getParametersBootstrap: %s Access connect client with bootstrap server.", LOG_LW2M_INFO, store.getEndPoint());
context.sendParametersOnThingsboardTelemetry(context.getKvLogyToThingsboard(logMsg), sessionInfo);
helper.sendParametersOnThingsboardTelemetry(helper.getKvLogyToThingsboard(logMsg), sessionInfo);
return lwM2MBootstrapConfig;
} else {
log.error(" [{}] Different values SecurityMode between of client and profile.", store.getEndPoint());
log.error("{} getParametersBootstrap: [{}] Different values SecurityMode between of client and profile.", LOG_LW2M_ERROR, store.getEndPoint());
String logMsg = String.format("%s: getParametersBootstrap: %s Different values SecurityMode between of client and profile.", LOG_LW2M_ERROR, store.getEndPoint());
context.sendParametersOnThingsboardTelemetry(context.getKvLogyToThingsboard(logMsg), sessionInfo);
helper.sendParametersOnThingsboardTelemetry(helper.getKvLogyToThingsboard(logMsg), sessionInfo);
return null;
}
}

17
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/secure/LwM2mCredentialsSecurityInfoValidator.java

@ -16,6 +16,7 @@
package org.thingsboard.server.transport.lwm2m.secure;
import com.google.gson.JsonObject;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.core.util.SecurityUtil;
@ -26,7 +27,9 @@ import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceLwM2MCredentialsRequestMsg;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContextServer;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandlerUtil;
import java.io.IOException;
@ -44,13 +47,11 @@ import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.X5
@Slf4j
@Component
@TbLwM2mTransportComponent
@RequiredArgsConstructor
public class LwM2mCredentialsSecurityInfoValidator {
private final LwM2mTransportContextServer contextS;
public LwM2mCredentialsSecurityInfoValidator(LwM2mTransportContextServer contextS) {
this.contextS = contextS;
}
private final LwM2mTransportContext context;
private final LwM2MTransportServerConfig config;
/**
* Request to thingsboard Response from thingsboard ValidateDeviceLwM2MCredentials
@ -61,7 +62,7 @@ public class LwM2mCredentialsSecurityInfoValidator {
public ReadResultSecurityStore createAndValidateCredentialsSecurityInfo(String endpoint, LwM2mTransportHandlerUtil.LwM2mTypeServer keyValue) {
CountDownLatch latch = new CountDownLatch(1);
final ReadResultSecurityStore[] resultSecurityStore = new ReadResultSecurityStore[1];
contextS.getTransportService().process(ValidateDeviceLwM2MCredentialsRequestMsg.newBuilder().setCredentialsId(endpoint).build(),
context.getTransportService().process(ValidateDeviceLwM2MCredentialsRequestMsg.newBuilder().setCredentialsId(endpoint).build(),
new TransportServiceCallback<>() {
@Override
public void onSuccess(ValidateDeviceCredentialsResponseMsg msg) {
@ -81,7 +82,7 @@ public class LwM2mCredentialsSecurityInfoValidator {
}
});
try {
latch.await(contextS.getLwM2MTransportServerConfig().getTimeout(), TimeUnit.MILLISECONDS);
latch.await(config.getTimeout(), TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
log.error("Failed to await credentials!", e);
}

105
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServiceImpl.java → common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2MTransportMsgHandler.java

@ -54,6 +54,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotif
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.client.LwM2mClient;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientProfile;
@ -111,44 +112,43 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandle
@Slf4j
@Service
@TbLwM2mTransportComponent
public class LwM2mTransportServiceImpl implements LwM2mTransportService {
public class DefaultLwM2MTransportMsgHandler implements LwM2mTransportMsgHandler {
private ExecutorService executorRegistered;
private ExecutorService executorUpdateRegistered;
private ExecutorService executorUnRegistered;
private LwM2mValueConverterImpl converter;
private FirmwareDataCache firmwareDataCache;
private final TransportService transportService;
public final LwM2mTransportContextServer lwM2mTransportContextServer;
private final LwM2mTransportContext context;
private final LwM2MTransportServerConfig config;
private final FirmwareDataCache firmwareDataCache;
private final LwM2mTransportServerHelper helper;
private final LwM2mClientContext lwM2mClientContext;
private final LeshanServer leshanServer;
private final LwM2mTransportRequest lwM2mTransportRequest;
public LwM2mTransportServiceImpl(TransportService transportService, LwM2mTransportContextServer lwM2mTransportContextServer,
LwM2mClientContext lwM2mClientContext, LeshanServer leshanServer,
@Lazy LwM2mTransportRequest lwM2mTransportRequest, FirmwareDataCache firmwareDataCache) {
public DefaultLwM2MTransportMsgHandler(TransportService transportService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper,
LwM2mClientContext lwM2mClientContext,
@Lazy LwM2mTransportRequest lwM2mTransportRequest,
FirmwareDataCache firmwareDataCache,
LwM2mTransportContext context) {
this.transportService = transportService;
this.lwM2mTransportContextServer = lwM2mTransportContextServer;
this.config = config;
this.helper = helper;
this.lwM2mClientContext = lwM2mClientContext;
this.leshanServer = leshanServer;
this.lwM2mTransportRequest = lwM2mTransportRequest;
this.firmwareDataCache = firmwareDataCache;
this.context = context;
}
@PostConstruct
public void init() {
this.lwM2mTransportContextServer.getScheduler().scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) lwM2mTransportContextServer.getLwM2MTransportServerConfig().getSessionReportTimeout()), lwM2mTransportContextServer.getLwM2MTransportServerConfig().getSessionReportTimeout(), TimeUnit.MILLISECONDS);
this.executorRegistered = Executors.newFixedThreadPool(this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getRegisteredPoolSize(),
this.context.getScheduler().scheduleAtFixedRate(this::checkInactivityAndReportActivity, new Random().nextInt((int) config.getSessionReportTimeout()), config.getSessionReportTimeout(), TimeUnit.MILLISECONDS);
this.executorRegistered = Executors.newFixedThreadPool(this.config.getRegisteredPoolSize(),
new NamedThreadFactory(String.format("LwM2M %s channel registered", SERVICE_CHANNEL)));
this.executorUpdateRegistered = Executors.newFixedThreadPool(this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getUpdateRegisteredPoolSize(),
this.executorUpdateRegistered = Executors.newFixedThreadPool(this.config.getUpdateRegisteredPoolSize(),
new NamedThreadFactory(String.format("LwM2M %s channel update registered", SERVICE_CHANNEL)));
this.executorUnRegistered = Executors.newFixedThreadPool(this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getUnRegisteredPoolSize(),
this.executorUnRegistered = Executors.newFixedThreadPool(this.config.getUnRegisteredPoolSize(),
new NamedThreadFactory(String.format("LwM2M %s channel un registered", SERVICE_CHANNEL)));
this.converter = LwM2mValueConverterImpl.getInstance();
}
@ -278,10 +278,10 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
@Override
public void setCancelObservations(Registration registration) {
if (registration != null) {
Set<Observation> observations = leshanServer.getObservationService().getObservations(registration);
Set<Observation> observations = context.getServer().getObservationService().getObservations(registration);
observations.forEach(observation -> lwM2mTransportRequest.sendAllRequest(registration,
convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), OBSERVE_CANCEL,
null, null, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null));
null, null, this.config.getTimeout(), null));
}
}
@ -333,13 +333,13 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
msg.getSharedUpdatedList().forEach(tsKvProto -> {
String pathName = tsKvProto.getKv().getKey();
String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, pathName);
Object valueNew = this.lwM2mTransportContextServer.getValueFromKvProto(tsKvProto.getKv());
Object valueNew = this.helper.getValueFromKvProto(tsKvProto.getKv());
//TODO: react on change of the firmware name.
if (FirmwareUtil.getAttributeKey(FirmwareType.FIRMWARE, FirmwareKey.VERSION).equals(pathName) && !valueNew.equals(lwM2MClient.getFrUpdate().getCurrentFwVersion())) {
this.getInfoFirmwareUpdate(lwM2MClient);
}
if (pathIdVer != null) {
ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig()
ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.config
.getModelProvider());
if (resourceModel != null && resourceModel.operations.isWritable()) {
this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), valueNew, pathIdVer);
@ -360,7 +360,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
} else if (msg.getSharedDeletedCount() > 0) {
msg.getSharedUpdatedList().forEach(tsKvProto -> {
String pathName = tsKvProto.getKv().getKey();
Object valueNew = this.lwM2mTransportContextServer.getValueFromKvProto(tsKvProto.getKv());
Object valueNew = this.helper.getValueFromKvProto(tsKvProto.getKv());
if (FirmwareUtil.getAttributeKey(FirmwareType.FIRMWARE, FirmwareKey.VERSION).equals(pathName) && !valueNew.equals(lwM2MClient.getFrUpdate().getCurrentFwVersion())) {
lwM2MClient.getFrUpdate().setCurrentFwVersion((String) valueNew);
}
@ -404,7 +404,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
@Override
public void onResourceUpdate(Optional<TransportProtos.ResourceUpdateMsg> resourceUpdateMsgOpt) {
String idVer = resourceUpdateMsgOpt.get().getResourceKey();
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.updateResourceModel(idVer, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getModelProvider()));
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.updateResourceModel(idVer, this.config.getModelProvider()));
}
/**
@ -413,7 +413,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
@Override
public void onResourceDelete(Optional<TransportProtos.ResourceDeleteMsg> resourceDeleteMsgOpt) {
String pathIdVer = resourceDeleteMsgOpt.get().getResourceKey();
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.deleteResources(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getModelProvider()));
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.deleteResources(pathIdVer, this.config.getModelProvider()));
}
@Override
@ -429,7 +429,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
} else {
lwM2mTransportRequest.sendAllRequest(registration, lwm2mClientRpcRequest.getTargetIdVer(), lwm2mClientRpcRequest.getTypeOper(), lwm2mClientRpcRequest.getContentFormatName(),
lwm2mClientRpcRequest.getValue() == null ? lwm2mClientRpcRequest.getParams() : lwm2mClientRpcRequest.getValue(),
this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), lwm2mClientRpcRequest);
this.config.getTimeout(), lwm2mClientRpcRequest);
}
} catch (Exception e) {
if (lwm2mClientRpcRequest == null) {
@ -546,7 +546,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
@Override
public void doTrigger(Registration registration, String path) {
lwM2mTransportRequest.sendAllRequest(registration, path, EXECUTE,
ContentFormat.TLV.getName(), null, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
ContentFormat.TLV.getName(), null, this.config.getTimeout(), null);
}
/**
@ -579,7 +579,8 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
*
* @param registration -
*/
protected void onAwakeDev(Registration registration) {
@Override
public void onAwakeDev(Registration registration) {
log.info("[{}] [{}] Received endpoint Awake version event", registration.getId(), registration.getEndpoint());
this.sendLogsToThingsboard(LOG_LW2M_INFO + ": Client is awake!", registration.getId());
//TODO: associate endpointId with device information.
@ -610,13 +611,14 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* @param logMsg - text msg
* @param registrationId - Id of Registration LwM2M Client
*/
@Override
public void sendLogsToThingsboard(String logMsg, String registrationId) {
SessionInfoProto sessionInfo = this.getValidateSessionInfo(registrationId);
if (logMsg != null && sessionInfo != null) {
if(logMsg.length() > 1024){
logMsg = logMsg.substring(0, 1024);
}
this.lwM2mTransportContextServer.sendParametersOnThingsboardTelemetry(this.lwM2mTransportContextServer.getKvLogyToThingsboard(logMsg), sessionInfo);
this.helper.sendParametersOnThingsboardTelemetry(this.helper.getKvLogyToThingsboard(logMsg), sessionInfo);
}
}
@ -640,7 +642,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
// #2
lwM2MClient.getPendingReadRequests().addAll(clientObjects);
clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(registration, path, READ, ContentFormat.TLV.getName(),
null, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null));
null, this.config.getTimeout(), null));
}
// #1
this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, READ, clientObjects);
@ -689,7 +691,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
*/
private void updateResourcesValue(Registration registration, LwM2mResource lwM2mResource, String path) {
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null);
if (lwM2MClient.saveResourceValue(path, lwM2mResource, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig()
if (lwM2MClient.saveResourceValue(path, lwM2mResource, this.config
.getModelProvider())) {
if (FR_PATH_RESOURCE_VER_ID.equals(convertPathFromIdVerToObjectId(path)) &&
lwM2MClient.getFrUpdate().getCurrentFwVersion() != null
@ -728,10 +730,10 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
SessionInfoProto sessionInfo = this.getValidateSessionInfo(registration);
if (results != null && sessionInfo != null) {
if (results.getResultAttributes().size() > 0) {
this.lwM2mTransportContextServer.sendParametersOnThingsboardAttribute(results.getResultAttributes(), sessionInfo);
this.helper.sendParametersOnThingsboardAttribute(results.getResultAttributes(), sessionInfo);
}
if (results.getResultTelemetries().size() > 0) {
this.lwM2mTransportContextServer.sendParametersOnThingsboardTelemetry(results.getResultTelemetries(), sessionInfo);
this.helper.sendParametersOnThingsboardTelemetry(results.getResultTelemetries(), sessionInfo);
}
}
} catch (Exception e) {
@ -780,7 +782,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
ConcurrentHashMap<String, Object> finalParams = params;
pathSend.forEach(target -> {
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, ContentFormat.TLV.getName(),
finalParams != null ? finalParams.get(target) : null, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
finalParams != null ? finalParams.get(target) : null, this.config.getTimeout(), null);
});
if (OBSERVE.equals(typeOper)) {
lwM2MClient.initReadValue(this, null);
@ -865,7 +867,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
LwM2mResource resourceValue = lwM2MClient != null ? getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer) : null;
if (resourceValue != null) {
ResourceModel.Type currentType = resourceValue.getType();
ResourceModel.Type expectedType = this.lwM2mTransportContextServer.getResourceModelTypeEqualsKvProtoValueType(currentType, pathIdVer);
ResourceModel.Type expectedType = this.helper.getResourceModelTypeEqualsKvProtoValueType(currentType, pathIdVer);
Object valueKvProto = null;
if (resourceValue.isMultiInstances()) {
valueKvProto = new JsonObject();
@ -882,7 +884,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
valueKvProto = this.converter.convertValue(resourceValue.getValue(), currentType, expectedType,
new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)));
}
return valueKvProto != null ? this.lwM2mTransportContextServer.getKvAttrTelemetryToThingsboard(currentType, resourceName, valueKvProto, resourceValue.isMultiInstances()) : null;
return valueKvProto != null ? this.helper.getKvAttrTelemetryToThingsboard(currentType, resourceName, valueKvProto, resourceValue.isMultiInstances()) : null;
}
} catch (Exception e) {
log.error("Failed to add parameters.", e);
@ -901,7 +903,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
private Object getResourceValueFormatKv(LwM2mClient lwM2MClient, String pathIdVer) {
LwM2mResource resourceValue = this.getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer);
ResourceModel.Type currentType = resourceValue.getType();
ResourceModel.Type expectedType = this.lwM2mTransportContextServer.getResourceModelTypeEqualsKvProtoValueType(currentType, pathIdVer);
ResourceModel.Type expectedType = this.helper.getResourceModelTypeEqualsKvProtoValueType(currentType, pathIdVer);
return this.converter.convertValue(resourceValue.getValue(), currentType, expectedType,
new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)));
}
@ -1094,10 +1096,10 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
if (pathIds.isResource()) {
if (READ.equals(typeOper)) {
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper,
ContentFormat.TLV.getName(), null, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
ContentFormat.TLV.getName(), null, this.config.getTimeout(), null);
} else if (OBSERVE.equals(typeOper)) {
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper,
null, null, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
null, null, this.config.getTimeout(), null);
}
}
});
@ -1153,7 +1155,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
if (!pathSend.isEmpty()) {
ConcurrentHashMap<String, Object> finalParams = lwm2mAttributesNew;
pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
finalParams.get(target), this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null));
finalParams.get(target), this.config.getTimeout(), null));
}
});
}
@ -1170,7 +1172,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
params.clear();
params.put(OBJECT_VERSION, "");
lwM2mTransportRequest.sendAllRequest(registration, target, WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
params, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
params, this.config.getTimeout(), null);
});
}
});
@ -1183,7 +1185,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
paramAnallyzer.forEach(pathIdVer -> {
if (this.getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer) != null) {
lwM2mTransportRequest.sendAllRequest(registration, pathIdVer, OBSERVE_CANCEL, null,
null, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
null, this.config.getTimeout(), null);
}
}
);
@ -1193,7 +1195,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
if (valueNew != null && (valueOld == null || !valueNew.toString().equals(valueOld.toString()))) {
lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, WRITE_REPLACE,
ContentFormat.TLV.getName(), valueNew,
this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
this.config.getTimeout(), null);
} else {
log.error("Failed update resource [{}] [{}]", path, valueNew);
String logMsg = String.format("%s: Failed update resource path - %s value - %s. Value is not changed or bad",
@ -1271,7 +1273,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
// #2.1
lwM2MClient.getDelayedRequests().forEach((pathIdVer, tsKvProto) -> {
this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer),
this.lwM2mTransportContextServer.getValueFromKvProto(tsKvProto.getKv()), pathIdVer);
this.helper.getValueFromKvProto(tsKvProto.getKv()), pathIdVer);
});
}
@ -1288,7 +1290,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
return null;
} else {
return SessionInfoProto.newBuilder()
.setNodeId(this.lwM2mTransportContextServer.getNodeId())
.setNodeId(this.context.getNodeId())
.setSessionIdMSB(lwM2MClient.getSessionId().getMostSignificantBits())
.setSessionIdLSB(lwM2MClient.getSessionId().getLeastSignificantBits())
.setDeviceIdMSB(msg.getDeviceInfo().getDeviceIdMSB())
@ -1358,7 +1360,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
if (keyNamesMap.values().size() > 0) {
try {
//#1.2
TransportProtos.GetAttributeRequestMsg getAttributeMsg = lwM2mTransportContextServer.getAdaptor().convertToGetAttributes(null, keyNamesMap.values());
TransportProtos.GetAttributeRequestMsg getAttributeMsg = helper.getAdaptor().convertToGetAttributes(null, keyNamesMap.values());
transportService.process(sessionInfo, getAttributeMsg, getAckCallback(lwM2MClient, getAttributeMsg.getRequestId(), DEVICE_ATTRIBUTES_REQUEST));
} catch (AdaptorException e) {
log.warn("Failed to decode get attributes request", e);
@ -1406,7 +1408,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
public void readRequestToClientFirmwareVer(Registration registration) {
String pathIdVer = convertPathFromObjectIdToIdVer(FR_PATH_RESOURCE_VER_ID, registration);
lwM2mTransportRequest.sendAllRequest(registration, pathIdVer, READ, ContentFormat.TLV.getName(),
null, lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
null, config.getTimeout(), null);
}
/**
@ -1422,7 +1424,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
String verSupportedObject = lwM2MClient.getRegistration().getSupportedObject().get(objectId);
String targetIdVer = LWM2M_SEPARATOR_PATH + objectId + LWM2M_SEPARATOR_KEY + verSupportedObject + LWM2M_SEPARATOR_PATH + 0 + LWM2M_SEPARATOR_PATH + 0;
lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), targetIdVer, WRITE_REPLACE, ContentFormat.OPAQUE.getName(),
firmwareChunk, lwM2mTransportContextServer.getLwM2MTransportServerConfig().getTimeout(), null);
firmwareChunk, config.getTimeout(), null);
log.warn("updateFirmwareClient [{}] [{}]", lwM2MClient.getFrUpdate().getCurrentFwVersion(), lwM2MClient.getFrUpdate().getClientFwVersion());
}
}
@ -1444,7 +1446,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
}
private boolean validateResourceInModel(LwM2mClient lwM2mClient, String pathIdVer, boolean isWritableNotOptional) {
ResourceModel resourceModel = lwM2mClient.getResourceModel(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig()
ResourceModel resourceModel = lwM2mClient.getResourceModel(pathIdVer, this.config
.getModelProvider());
Integer objectId = new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)).getObjectId();
String objectVer = validateObjectVerFromKey(pathIdVer);
@ -1453,9 +1455,4 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
objectId != null && objectVer != null && objectVer.equals(lwM2mClient.getRegistration().getSupportedVersion(objectId)));
}
@Override
public String getName() {
return "LWM2M";
}
}

77
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerConfiguration.java → common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.transport.lwm2m.server;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.californium.core.network.config.NetworkConfig;
import org.eclipse.californium.core.network.stack.BlockwiseLayer;
@ -29,14 +30,16 @@ import org.eclipse.leshan.server.model.LwM2mModelProvider;
import org.eclipse.leshan.server.security.DefaultAuthorizer;
import org.eclipse.leshan.server.security.EditableSecurityStore;
import org.eclipse.leshan.server.security.SecurityChecker;
import org.springframework.context.annotation.Bean;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import org.thingsboard.server.transport.lwm2m.secure.LWM2MGenerationPSkRPkECC;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
import javax.annotation.PostConstruct;
import javax.annotation.PreDestroy;
import java.math.BigInteger;
import java.security.AlgorithmParameters;
import java.security.KeyFactory;
@ -66,32 +69,53 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mNetworkConfig.g
@Slf4j
@Component
@TbLwM2mTransportComponent
public class LwM2mTransportServerConfiguration {
@RequiredArgsConstructor
public class DefaultLwM2mTransportService implements LwM2MTransportService {
private PublicKey publicKey;
private PrivateKey privateKey;
private boolean pskMode = false;
private final LwM2mTransportContextServer context;
private final LwM2mTransportContext context;
private final LwM2MTransportServerConfig config;
private final LwM2mTransportServerHelper helper;
private final LwM2mTransportMsgHandler handler;
private final CaliforniumRegistrationStore registrationStore;
private final EditableSecurityStore securityStore;
private final LwM2mClientContext lwM2mClientContext;
public LwM2mTransportServerConfiguration(LwM2mTransportContextServer context, CaliforniumRegistrationStore registrationStore, EditableSecurityStore securityStore, LwM2mClientContext lwM2mClientContext) {
this.context = context;
this.registrationStore = registrationStore;
this.securityStore = securityStore;
this.lwM2mClientContext = lwM2mClientContext;
private LeshanServer server;
@PostConstruct
public void init() {
if (config.getEnableGenNewKeyPskRpk()) {
new LWM2MGenerationPSkRPkECC();
}
this.server = getLhServer(config.getPort(), config.getSecurePort());
this.startLhServer();
this.context.setServer(server);
}
private void startLhServer() {
log.info("Starting LwM2M transport Server...");
this.server.start();
LwM2mServerListener lhServerCertListener = new LwM2mServerListener(handler);
this.server.getRegistrationService().addListener(lhServerCertListener.registrationListener);
this.server.getPresenceService().addListener(lhServerCertListener.presenceListener);
this.server.getObservationService().addListener(lhServerCertListener.observationListener);
}
@Bean
public LeshanServer getLeshanServer() {
log.info("Starting LwM2M transport Server... PostConstruct");
return this.getLhServer(this.context.getLwM2MTransportServerConfig().getPort(), this.context.getLwM2MTransportServerConfig().getSecurePort());
@PreDestroy
public void shutdown() {
log.info("Stopping LwM2M transport Server!");
server.destroy();
log.info("LwM2M transport Server stopped!");
}
private LeshanServer getLhServer(Integer serverPortNoSec, Integer serverSecurePort) {
LeshanServerBuilder builder = new LeshanServerBuilder();
builder.setLocalAddress(this.context.getLwM2MTransportServerConfig().getHost(), serverPortNoSec);
builder.setLocalSecureAddress(this.context.getLwM2MTransportServerConfig().getSecureHost(), serverSecurePort);
builder.setLocalAddress(config.getHost(), serverPortNoSec);
builder.setLocalSecureAddress(config.getSecureHost(), serverSecurePort);
builder.setDecoder(new DefaultLwM2mNodeDecoder());
/** Use a magic converter to support bad type send by the UI. */
builder.setEncoder(new DefaultLwM2mNodeEncoder(LwM2mValueConverterImpl.getInstance()));
@ -103,8 +127,8 @@ public class LwM2mTransportServerConfiguration {
builder.setCoapConfig(getCoapConfig(serverPortNoSec, serverSecurePort));
/** Define model provider (Create Models )*/
LwM2mModelProvider modelProvider = new LwM2mVersionedModelProvider(this.lwM2mClientContext, this.context);
this.context.getLwM2MTransportServerConfig().setModelProvider(modelProvider);
LwM2mModelProvider modelProvider = new LwM2mVersionedModelProvider(this.lwM2mClientContext, this.helper, this.context);
config.setModelProvider(modelProvider);
builder.setObjectModelProvider(modelProvider);
/** Create credentials */
@ -118,8 +142,8 @@ public class LwM2mTransportServerConfiguration {
/** Create DTLS Config */
DtlsConnectorConfig.Builder dtlsConfig = new DtlsConnectorConfig.Builder();
dtlsConfig.setServerOnly(true);
dtlsConfig.setRecommendedSupportedGroupsOnly(this.context.getLwM2MTransportServerConfig().isRecommendedSupportedGroups());
dtlsConfig.setRecommendedCipherSuitesOnly(this.context.getLwM2MTransportServerConfig().isRecommendedCiphers());
dtlsConfig.setRecommendedSupportedGroupsOnly(config.isRecommendedSupportedGroups());
dtlsConfig.setRecommendedCipherSuitesOnly(config.isRecommendedCiphers());
if (this.pskMode) {
dtlsConfig.setSupportedCipherSuites(
TLS_PSK_WITH_AES_128_CCM_8,
@ -141,9 +165,9 @@ public class LwM2mTransportServerConfiguration {
private void setServerWithCredentials(LeshanServerBuilder builder) {
try {
if (this.context.getLwM2MTransportServerConfig().getKeyStoreValue() != null) {
if (config.getKeyStoreValue() != null) {
if (this.setBuilderX509(builder)) {
X509Certificate rootCAX509Cert = (X509Certificate) this.context.getLwM2MTransportServerConfig().getKeyStoreValue().getCertificate(this.context.getLwM2MTransportServerConfig().getRootCertificateAlias());
X509Certificate rootCAX509Cert = (X509Certificate) config.getKeyStoreValue().getCertificate(config.getRootCertificateAlias());
if (rootCAX509Cert != null) {
X509Certificate[] trustedCertificates = new X509Certificate[1];
trustedCertificates[0] = rootCAX509Cert;
@ -178,8 +202,8 @@ public class LwM2mTransportServerConfiguration {
private boolean setBuilderX509(LeshanServerBuilder builder) {
try {
X509Certificate serverCertificate = (X509Certificate) this.context.getLwM2MTransportServerConfig().getKeyStoreValue().getCertificate(this.context.getLwM2MTransportServerConfig().getCertificateAlias());
PrivateKey privateKey = (PrivateKey) this.context.getLwM2MTransportServerConfig().getKeyStoreValue().getKey(this.context.getLwM2MTransportServerConfig().getCertificateAlias(), this.context.getLwM2MTransportServerConfig().getKeyStorePassword() == null ? null : this.context.getLwM2MTransportServerConfig().getKeyStorePassword().toCharArray());
X509Certificate serverCertificate = (X509Certificate) config.getKeyStoreValue().getCertificate(config.getCertificateAlias());
PrivateKey privateKey = (PrivateKey) config.getKeyStoreValue().getKey(config.getCertificateAlias(), config.getKeyStorePassword() == null ? null : config.getKeyStorePassword().toCharArray());
PublicKey publicKey = serverCertificate.getPublicKey();
if (privateKey != null && privateKey.getEncoded().length > 0 && publicKey != null && publicKey.getEncoded().length > 0) {
builder.setPublicKey(serverCertificate.getPublicKey());
@ -208,7 +232,7 @@ public class LwM2mTransportServerConfiguration {
}
private void infoPramsUri(String mode) {
LwM2MTransportServerConfig lwM2MTransportServerConfig = this.context.getLwM2MTransportServerConfig();
LwM2MTransportServerConfig lwM2MTransportServerConfig = config;
log.info("Server uses [{}]: serverNoSecureURI : [{}:{}], serverSecureURI : [{}:{}]", mode,
lwM2MTransportServerConfig.getHost(),
lwM2MTransportServerConfig.getPort(),
@ -236,7 +260,7 @@ public class LwM2mTransportServerConfiguration {
AlgorithmParameters algoParameters = AlgorithmParameters.getInstance("EC");
algoParameters.init(new ECGenParameterSpec("secp256r1"));
ECParameterSpec parameterSpec = algoParameters.getParameterSpec(ECParameterSpec.class);
LwM2MTransportServerConfig serverConfig = this.context.getLwM2MTransportServerConfig();
LwM2MTransportServerConfig serverConfig = config;
if (StringUtils.isNotEmpty(serverConfig.getPublicX()) && StringUtils.isNotEmpty(serverConfig.getPublicY())) {
byte[] publicX = Hex.decodeHex(serverConfig.getPublicX().toCharArray());
byte[] publicY = Hex.decodeHex(serverConfig.getPublicY().toCharArray());
@ -283,4 +307,9 @@ public class LwM2mTransportServerConfiguration {
params);
}
@Override
public String getName() {
return "LWM2M";
}
}

22
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2MTransportService.java

@ -0,0 +1,22 @@
/**
* Copyright © 2016-2021 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* 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
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.server;
import org.thingsboard.server.common.data.TbTransportService;
public interface LwM2MTransportService extends TbTransportService {
}

4
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java

@ -32,9 +32,9 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandle
@Slf4j
public class LwM2mServerListener {
private final LwM2mTransportServiceImpl service;
private final LwM2mTransportMsgHandler service;
public LwM2mServerListener(LwM2mTransportServiceImpl service) {
public LwM2mServerListener(LwM2mTransportMsgHandler service) {
this.service = service;
}

4
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java

@ -35,10 +35,10 @@ import java.util.Optional;
@Slf4j
public class LwM2mSessionMsgListener implements GenericFutureListener<Future<? super Void>>, SessionMsgListener {
private LwM2mTransportServiceImpl service;
private DefaultLwM2MTransportMsgHandler service;
private TransportProtos.SessionInfoProto sessionInfo;
public LwM2mSessionMsgListener(LwM2mTransportServiceImpl service, TransportProtos.SessionInfoProto sessionInfo) {
public LwM2mSessionMsgListener(DefaultLwM2MTransportMsgHandler service, TransportProtos.SessionInfoProto sessionInfo) {
this.service = service;
this.sessionInfo = sessionInfo;
}

32
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportContext.java

@ -0,0 +1,32 @@
/**
* Copyright © 2016-2021 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* 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
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.server;
import lombok.Getter;
import lombok.Setter;
import org.eclipse.leshan.server.californium.LeshanServer;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.transport.TransportContext;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
@Component
@TbLwM2mTransportComponent
public class LwM2mTransportContext extends TransportContext {
@Getter @Setter
private LeshanServer server;
}

6
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportService.java → common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportMsgHandler.java

@ -27,7 +27,7 @@ import org.thingsboard.server.transport.lwm2m.server.client.Lwm2mClientRpcReques
import java.util.Collection;
import java.util.Optional;
public interface LwM2mTransportService extends TbTransportService {
public interface LwM2mTransportMsgHandler {
void onRegistered(Registration registration, Collection<Observation> previousObsersations);
@ -60,4 +60,8 @@ public interface LwM2mTransportService extends TbTransportService {
void doTrigger(Registration registration, String path);
void doDisconnect(TransportProtos.SessionInfoProto sessionInfo);
void onAwakeDev(Registration registration);
void sendLogsToThingsboard(String msg, String registrationId);
}

37
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportRequest.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.transport.lwm2m.server;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.californium.core.coap.CoAP;
@ -50,6 +51,7 @@ import org.eclipse.leshan.server.registration.Registration;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
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.Lwm2mClientRpcRequest;
@ -82,35 +84,22 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandle
@Slf4j
@Service
@TbLwM2mTransportComponent
@RequiredArgsConstructor
public class LwM2mTransportRequest {
private ExecutorService executorResponse;
public LwM2mValueConverterImpl converter;
private final LwM2mTransportContextServer lwM2mTransportContextServer;
private final LwM2mTransportContext context;
private final LwM2MTransportServerConfig config;
private final LwM2mTransportServerHelper lwM2MTransportServerHelper;
private final LwM2mClientContext lwM2mClientContext;
private final LeshanServer leshanServer;
private final LwM2mTransportServiceImpl serviceImpl;
private final TransportService transportService;
public LwM2mTransportRequest(LwM2mTransportContextServer lwM2mTransportContextServer,
LwM2mClientContext lwM2mClientContext, LeshanServer leshanServer,
LwM2mTransportServiceImpl serviceImpl, TransportService transportService) {
this.lwM2mTransportContextServer = lwM2mTransportContextServer;
this.lwM2mClientContext = lwM2mClientContext;
this.leshanServer = leshanServer;
this.serviceImpl = serviceImpl;
this.transportService = transportService;
}
private final DefaultLwM2MTransportMsgHandler serviceImpl;
@PostConstruct
public void init() {
this.converter = LwM2mValueConverterImpl.getInstance();
executorResponse = Executors.newFixedThreadPool(this.lwM2mTransportContextServer.getLwM2MTransportServerConfig().getResponsePoolSize(),
executorResponse = Executors.newFixedThreadPool(this.config.getResponsePoolSize(),
new NamedThreadFactory(String.format("LwM2M %s channel response", RESPONSE_CHANNEL)));
}
@ -158,10 +147,10 @@ public class LwM2mTransportRequest {
* At server side this will not remove the observation from the observation store, to do it you need to use
* {@code ObservationService#cancelObservation()}
*/
leshanServer.getObservationService().cancelObservations(registration, target);
context.getServer().getObservationService().cancelObservations(registration, target);
break;
case EXECUTE:
resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig()
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));
@ -171,7 +160,7 @@ public class LwM2mTransportRequest {
break;
case WRITE_REPLACE:
// Request to write a <b>String Single-Instance Resource</b> using the TLV content format.
resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.lwM2mTransportContextServer.getLwM2MTransportServerConfig()
resourceModel = lwM2MClient.getResourceModel(targetIdVer, this.config
.getModelProvider());
if (contentFormat.equals(ContentFormat.TLV)) {
request = this.getWriteRequestSingleResource(null, resultIds.getObjectId(),
@ -232,7 +221,7 @@ public class LwM2mTransportRequest {
serviceImpl.sentRpcRequest(rpcRequest, NOT_FOUND.getName(), errorMsg, LOG_LW2M_ERROR);
}
} else if (OBSERVE_READ_ALL.name().equals(typeOper.name())) {
Set<Observation> observations = leshanServer.getObservationService().getObservations(registration);
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);
@ -259,7 +248,7 @@ public class LwM2mTransportRequest {
@SuppressWarnings("unchecked")
private void sendRequest(Registration registration, LwM2mClient lwM2MClient, DownlinkRequest request, long timeoutInMs, Lwm2mClientRpcRequest rpcRequest) {
leshanServer.send(registration, request, timeoutInMs, (ResponseCallback<?>) response -> {
context.getServer().send(registration, request, timeoutInMs, (ResponseCallback<?>) response -> {
if (!lwM2MClient.isInit()) {
lwM2MClient.initReadValue(this.serviceImpl, convertPathFromObjectIdToIdVer(request.getPath().toString(), registration));
}

83
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportContextServer.java → common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerHelper.java

@ -31,6 +31,7 @@ package org.thingsboard.server.transport.lwm2m.server;
*/
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.model.DDFFileParser;
import org.eclipse.leshan.core.model.DefaultDDFFileValidator;
@ -39,11 +40,8 @@ import org.eclipse.leshan.core.model.ObjectModel;
import org.eclipse.leshan.core.model.ResourceModel;
import org.eclipse.leshan.core.node.codec.CodecException;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.transport.TransportContext;
import org.thingsboard.server.common.transport.TransportResourceCache;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.PostAttributeMsg;
import org.thingsboard.server.gen.transport.TransportProtos.PostTelemetryMsg;
@ -62,34 +60,16 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandle
@Slf4j
@Component
@TbLwM2mTransportComponent
public class LwM2mTransportContextServer extends TransportContext {
@RequiredArgsConstructor
public class LwM2mTransportServerHelper {
private final LwM2MTransportServerConfig lwM2MTransportServerConfig;
private final LwM2mTransportContext context;
private final TransportService transportService;
private final TransportResourceCache transportResourceCache;
@Getter
private final LwM2MJsonAdaptor adaptor;
public LwM2mTransportContextServer(LwM2MTransportServerConfig lwM2MTransportServerConfig, TransportService transportService, TransportResourceCache transportResourceCache, LwM2MJsonAdaptor adaptor) {
this.lwM2MTransportServerConfig = lwM2MTransportServerConfig;
this.transportService = transportService;
this.transportResourceCache = transportResourceCache;
this.adaptor = adaptor;
}
public LwM2MTransportServerConfig getLwM2MTransportServerConfig() {
return this.lwM2MTransportServerConfig;
}
public TransportResourceCache getTransportResourceCache() {
return this.transportResourceCache;
}
/**
* send to Thingsboard Attribute || Telemetry
*
@ -134,7 +114,7 @@ public class LwM2mTransportContextServer extends TransportContext {
*/
public SessionInfoProto getValidateSessionInfo(TransportProtos.ValidateDeviceCredentialsResponseMsg msg, long mostSignificantBits, long leastSignificantBits) {
return SessionInfoProto.newBuilder()
.setNodeId(this.getNodeId())
.setNodeId(context.getNodeId())
.setSessionIdMSB(mostSignificantBits)
.setSessionIdLSB(leastSignificantBits)
.setDeviceIdMSB(msg.getDeviceInfo().getDeviceIdMSB())
@ -165,8 +145,8 @@ public class LwM2mTransportContextServer extends TransportContext {
* @param logMsg - info about Logs
* @return- KeyValueProto for telemetry (Logs)
*/
public List <TransportProtos.KeyValueProto> getKvLogyToThingsboard(String logMsg) {
List <TransportProtos.KeyValueProto> result = new ArrayList<>();
public List<TransportProtos.KeyValueProto> getKvLogyToThingsboard(String logMsg) {
List<TransportProtos.KeyValueProto> result = new ArrayList<>();
result.add(TransportProtos.KeyValueProto.newBuilder()
.setKey(LOG_LW2M_TELEMETRY)
.setType(TransportProtos.KeyValueType.STRING_V)
@ -179,32 +159,31 @@ public class LwM2mTransportContextServer extends TransportContext {
* @throws CodecException -
*/
public TransportProtos.KeyValueProto getKvAttrTelemetryToThingsboard(ResourceModel.Type resourceType, String resourceName, Object value, boolean isMultiInstances) {
TransportProtos.KeyValueProto.Builder kvProto = TransportProtos.KeyValueProto.newBuilder().setKey(resourceName);
if (isMultiInstances) {
kvProto.setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV((String) value);
public TransportProtos.KeyValueProto getKvAttrTelemetryToThingsboard(ResourceModel.Type resourceType, String resourceName, Object value, boolean isMultiInstances) {
TransportProtos.KeyValueProto.Builder kvProto = TransportProtos.KeyValueProto.newBuilder().setKey(resourceName);
if (isMultiInstances) {
kvProto.setType(TransportProtos.KeyValueType.JSON_V)
.setJsonV((String) value);
} else {
switch (resourceType) {
case BOOLEAN:
kvProto.setType(BOOLEAN_V).setBoolV((Boolean) value).build();
break;
case STRING:
case TIME:
case OPAQUE:
case OBJLNK:
kvProto.setType(TransportProtos.KeyValueType.STRING_V).setStringV((String) value);
break;
case INTEGER:
kvProto.setType(TransportProtos.KeyValueType.LONG_V).setLongV((Long) value);
break;
case FLOAT:
kvProto.setType(TransportProtos.KeyValueType.DOUBLE_V).setDoubleV((Double) value);
}
else {
switch (resourceType) {
case BOOLEAN:
kvProto.setType(BOOLEAN_V).setBoolV((Boolean) value).build();
break;
case STRING:
case TIME:
case OPAQUE:
case OBJLNK:
kvProto.setType(TransportProtos.KeyValueType.STRING_V).setStringV((String) value);
break;
case INTEGER:
kvProto.setType(TransportProtos.KeyValueType.LONG_V).setLongV((Long) value);
break;
case FLOAT:
kvProto.setType(TransportProtos.KeyValueType.DOUBLE_V).setDoubleV((Double) value);
}
}
return kvProto.build();
}
return kvProto.build();
}
/**
*
@ -230,7 +209,7 @@ public class LwM2mTransportContextServer extends TransportContext {
throw new CodecException("Invalid ResourceModel_Type for resource %s, got %s", resourcePath, currentType);
}
public Object getValueFromKvProto (TransportProtos.KeyValueProto kv) {
public Object getValueFromKvProto(TransportProtos.KeyValueProto kv) {
switch (kv.getType()) {
case BOOLEAN_V:
return kv.getBoolV();

31
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportServerInitializer.java

@ -30,35 +30,4 @@ import javax.annotation.PreDestroy;
@TbLwM2mTransportComponent
public class LwM2mTransportServerInitializer {
@Autowired
private LwM2mTransportServiceImpl service;
@Autowired
private LeshanServer leshanServer;
@Autowired
private LwM2mTransportContextServer context;
@PostConstruct
public void init() {
if (this.context.getLwM2MTransportServerConfig().getEnableGenNewKeyPskRpk()) {
new LWM2MGenerationPSkRPkECC();
}
this.startLhServer();
}
private void startLhServer() {
this.leshanServer.start();
LwM2mServerListener lhServerCertListener = new LwM2mServerListener(service);
this.leshanServer.getRegistrationService().addListener(lhServerCertListener.registrationListener);
this.leshanServer.getPresenceService().addListener(lhServerCertListener.presenceListener);
this.leshanServer.getObservationService().addListener(lhServerCertListener.observationListener);
}
@PreDestroy
public void shutdown() {
log.info("Stopping LwM2M transport Server!");
leshanServer.destroy();
log.info("LwM2M transport Server stopped!");
}
}

16
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mVersionedModelProvider.java

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.transport.lwm2m.server;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.model.DefaultDDFFileValidator;
import org.eclipse.leshan.core.model.LwM2mModel;
@ -37,6 +38,7 @@ import static org.thingsboard.server.common.data.ResourceType.LWM2M_MODEL;
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_KEY;
@Slf4j
@RequiredArgsConstructor
public class LwM2mVersionedModelProvider implements LwM2mModelProvider {
/**
@ -46,12 +48,8 @@ public class LwM2mVersionedModelProvider implements LwM2mModelProvider {
* Value = TenantId
*/
private final LwM2mClientContext lwM2mClientContext;
private final LwM2mTransportContextServer lwM2mTransportContextServer;
public LwM2mVersionedModelProvider(LwM2mClientContext lwM2mClientContext, LwM2mTransportContextServer lwM2mTransportContextServer) {
this.lwM2mClientContext = lwM2mClientContext;
this.lwM2mTransportContextServer = lwM2mTransportContextServer;
}
private final LwM2mTransportServerHelper helper;
private final LwM2mTransportContext context;
private String getKeyIdVer(Integer objectId, String version) {
return objectId != null ? objectId + LWM2M_SEPARATOR_KEY + ((version == null || version.isEmpty()) ? ObjectModel.DEFAULT_VERSION : version) : null;
@ -120,11 +118,9 @@ public class LwM2mVersionedModelProvider implements LwM2mModelProvider {
private ObjectModel getObjectModelDynamic(Integer objectId, String version) {
String key = getKeyIdVer(objectId, version);
Optional<TbResource> tbResource = lwM2mTransportContextServer
.getTransportResourceCache()
.get(this.tenantId, LWM2M_MODEL, key);
Optional<TbResource> tbResource = context.getTransportResourceCache().get(this.tenantId, LWM2M_MODEL, key);
return tbResource.map(resource -> lwM2mTransportContextServer.parseFromXmlToObjectModel(
return tbResource.map(resource -> helper.parseFromXmlToObjectModel(
Base64.getDecoder().decode(resource.getData()),
key + ".xml",
new DefaultDDFFileValidator())).orElse(null);

4
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java

@ -27,7 +27,7 @@ import org.eclipse.leshan.server.security.SecurityInfo;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServiceImpl;
import org.thingsboard.server.transport.lwm2m.server.DefaultLwM2MTransportMsgHandler;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
import java.util.Collection;
@ -172,7 +172,7 @@ public class LwM2mClient implements Cloneable {
.collect(Collectors.toSet());
}
public void initReadValue(LwM2mTransportServiceImpl serviceImpl, String path) {
public void initReadValue(DefaultLwM2MTransportMsgHandler serviceImpl, String path) {
if (path != null) {
this.pendingReadRequests.remove(path);
}

4
common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/TransportContext.java

@ -51,11 +51,13 @@ public abstract class TransportContext {
@Getter
private ExecutorService executor;
@Getter
@Autowired
private FirmwareDataCache firmwareDataCache;
@Autowired
private TransportResourceCache transportResourceCache;
@PostConstruct
public void init() {
executor = ThingsBoardExecutors.newWorkStealingPool(50, getClass());

Loading…
Cancel
Save