From bd35adc5ba9ae65cae06b29d551c45b6f863d736 Mon Sep 17 00:00:00 2001 From: Samuel Gallet Date: Thu, 3 Mar 2022 10:40:38 +0100 Subject: [PATCH 1/2] Implement SendService Listener --- .../server/DefaultLwM2mTransportService.java | 6 +++ .../lwm2m/server/LwM2mServerListener.java | 14 ++++++ .../uplink/DefaultLwM2mUplinkMsgHandler.java | 48 +++++++++++++------ .../server/uplink/LwM2mUplinkMsgHandler.java | 3 ++ 4 files changed, 56 insertions(+), 15 deletions(-) diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java index 7860fd94a5..9b49030f7e 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/DefaultLwM2mTransportService.java @@ -20,11 +20,15 @@ import lombok.extern.slf4j.Slf4j; import org.eclipse.californium.scandium.config.DtlsConfig; import org.eclipse.californium.scandium.config.DtlsConnectorConfig; import org.eclipse.californium.scandium.dtls.cipher.CipherSuite; +import org.eclipse.leshan.core.node.LwM2mNode; import org.eclipse.leshan.core.node.codec.DefaultLwM2mDecoder; import org.eclipse.leshan.core.node.codec.DefaultLwM2mEncoder; +import org.eclipse.leshan.core.request.SendRequest; import org.eclipse.leshan.server.californium.LeshanServer; import org.eclipse.leshan.server.californium.LeshanServerBuilder; import org.eclipse.leshan.server.californium.registration.CaliforniumRegistrationStore; +import org.eclipse.leshan.server.registration.Registration; +import org.eclipse.leshan.server.send.SendListener; import org.springframework.stereotype.Component; import org.thingsboard.server.cache.ota.OtaPackageDataCache; import org.thingsboard.server.common.data.DataConstants; @@ -40,6 +44,7 @@ import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PreDestroy; import java.security.cert.X509Certificate; +import java.util.Map; import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_RECOMMENDED_CIPHER_SUITES_ONLY; import static org.eclipse.californium.scandium.config.DtlsConfig.DTLS_RECOMMENDED_CURVES_ONLY; @@ -94,6 +99,7 @@ public class DefaultLwM2mTransportService implements LwM2MTransportService { this.server.getRegistrationService().addListener(lhServerCertListener.registrationListener); this.server.getPresenceService().addListener(lhServerCertListener.presenceListener); this.server.getObservationService().addListener(lhServerCertListener.observationListener); + this.server.getSendService().addListener(lhServerCertListener.sendListener); log.info("Started LwM2M transport server."); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java index 5c0e225f8b..4d4af68fad 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java @@ -16,9 +16,11 @@ package org.thingsboard.server.transport.lwm2m.server; import lombok.extern.slf4j.Slf4j; +import org.eclipse.leshan.core.node.LwM2mNode; import org.eclipse.leshan.core.observation.CompositeObservation; import org.eclipse.leshan.core.observation.Observation; import org.eclipse.leshan.core.observation.SingleObservation; +import org.eclipse.leshan.core.request.SendRequest; import org.eclipse.leshan.core.response.ObserveCompositeResponse; import org.eclipse.leshan.core.response.ObserveResponse; import org.eclipse.leshan.server.observation.ObservationListener; @@ -26,9 +28,11 @@ import org.eclipse.leshan.server.queue.PresenceListener; import org.eclipse.leshan.server.registration.Registration; import org.eclipse.leshan.server.registration.RegistrationListener; import org.eclipse.leshan.server.registration.RegistrationUpdate; +import org.eclipse.leshan.server.send.SendListener; import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import java.util.Collection; +import java.util.Map; import static org.thingsboard.server.transport.lwm2m.utils.LwM2MTransportUtil.convertObjectIdToVersionedId; @@ -119,4 +123,14 @@ public class LwM2mServerListener { log.trace("Successful start newObservation {}.", ((SingleObservation)observation).getPath()); } }; + + public final SendListener sendListener = new SendListener() { + + @Override + public void dataReceived(Registration registration, Map map, SendRequest sendRequest) { + if (registration != null) { + service.onUpdateValueWithSendRequest(registration, sendRequest); + } + } + }; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java index 05ce126c7e..0f11bd47df 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java @@ -33,11 +33,7 @@ import org.eclipse.leshan.core.node.LwM2mResource; import org.eclipse.leshan.core.node.LwM2mResourceInstance; import org.eclipse.leshan.core.node.LwM2mSingleResource; import org.eclipse.leshan.core.observation.Observation; -import org.eclipse.leshan.core.request.CreateRequest; -import org.eclipse.leshan.core.request.ObserveRequest; -import org.eclipse.leshan.core.request.ReadRequest; -import org.eclipse.leshan.core.request.WriteCompositeRequest; -import org.eclipse.leshan.core.request.WriteRequest; +import org.eclipse.leshan.core.request.*; import org.eclipse.leshan.core.request.WriteRequest.Mode; import org.eclipse.leshan.core.response.ObserveResponse; import org.eclipse.leshan.core.response.ReadCompositeResponse; @@ -98,16 +94,7 @@ import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; -import java.util.ArrayList; -import java.util.Collection; -import java.util.Collections; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.Optional; -import java.util.Random; -import java.util.Set; -import java.util.UUID; +import java.util.*; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -382,6 +369,37 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl } } + @Override + public void onUpdateValueWithSendRequest(Registration registration, SendRequest sendRequest) { + Iterator i$ = sendRequest.getNodes().entrySet().iterator(); + if (i$.hasNext()) { + Map.Entry entry = (Map.Entry) i$.next(); + LwM2mPath path = (LwM2mPath) entry.getKey(); + LwM2mNode node = (LwM2mNode) entry.getValue(); + LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint()); + String stringPath = convertObjectIdToVersionedId(path.toString(), registration); + ObjectModel objectModelVersion = lwM2MClient.getObjectModel(stringPath, modelProvider); + if (objectModelVersion != null) { + if (node instanceof LwM2mObject) { + LwM2mObject lwM2mObject = (LwM2mObject) node; + this.updateObjectResourceValue(lwM2MClient, lwM2mObject, stringPath, 0); + } else if (node instanceof LwM2mObjectInstance) { + LwM2mObjectInstance lwM2mObjectInstance = (LwM2mObjectInstance) node; + this.updateObjectInstanceResourceValue(lwM2MClient, lwM2mObjectInstance, stringPath, 0); + } else if (node instanceof LwM2mResource) { + LwM2mResource lwM2mResource = (LwM2mResource) node; + this.updateResourcesValue(lwM2MClient, lwM2mResource, stringPath, Mode.UPDATE, 0); + } + } + if (clientContext.awake(lwM2MClient)) { + // clientContext.awake calls clientContext.update + log.debug("[{}] Device is awake", lwM2MClient.getEndpoint()); + } else { + clientContext.update(lwM2MClient); + } + } + } + /** * @param sessionInfo - * @param deviceProfile - diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java index 475d834c60..332bb6ccad 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java @@ -17,6 +17,7 @@ package org.thingsboard.server.transport.lwm2m.server.uplink; import org.eclipse.leshan.core.observation.Observation; import org.eclipse.leshan.core.request.CreateRequest; +import org.eclipse.leshan.core.request.SendRequest; import org.eclipse.leshan.core.request.WriteCompositeRequest; import org.eclipse.leshan.core.request.WriteRequest; import org.eclipse.leshan.core.response.ReadCompositeResponse; @@ -46,6 +47,8 @@ public interface LwM2mUplinkMsgHandler { void onUpdateValueAfterReadCompositeResponse(Registration registration, ReadCompositeResponse response); + void onUpdateValueWithSendRequest(Registration registration, SendRequest sendRequest); + void onDeviceProfileUpdate(TransportProtos.SessionInfoProto sessionInfo, DeviceProfile deviceProfile); void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional deviceProfileOpt); From 51c88b7712a54a8476648eaff86b76603f919497 Mon Sep 17 00:00:00 2001 From: Samuel Gallet Date: Thu, 3 Mar 2022 10:59:35 +0100 Subject: [PATCH 2/2] Add method documentation --- .../lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java index 0f11bd47df..1dcd4cbab6 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2mUplinkMsgHandler.java @@ -369,6 +369,12 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl } } + /** + * Sending updated value to thingsboard from SendListener.dataReceived: object, instance, SingleResource or MultipleResource + * + * @param registration - Registration LwM2M Client + * @param sendRequest - sendRequest + */ @Override public void onUpdateValueWithSendRequest(Registration registration, SendRequest sendRequest) { Iterator i$ = sendRequest.getNodes().entrySet().iterator();