|
|
|
@ -60,6 +60,7 @@ import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTrans |
|
|
|
import org.thingsboard.server.common.data.device.profile.lwm2m.ObjectAttributes; |
|
|
|
import org.thingsboard.server.common.data.device.profile.lwm2m.OtherConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryMappingConfiguration; |
|
|
|
import org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.ota.OtaPackageUtil; |
|
|
|
@ -92,6 +93,8 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadCallbac |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MReadRequest; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesCallback; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttributesRequest; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MObserveCompositeCallback; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MObserveCompositeRequest; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.model.LwM2MModelConfig; |
|
|
|
import org.thingsboard.server.transport.lwm2m.server.model.LwM2MModelConfigService; |
|
|
|
@ -116,9 +119,11 @@ import java.util.UUID; |
|
|
|
import java.util.concurrent.ConcurrentHashMap; |
|
|
|
import java.util.concurrent.CountDownLatch; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.concurrent.atomic.AtomicLong; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_ALL; |
|
|
|
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.COMPOSITE_BY_OBJECT; |
|
|
|
import static org.thingsboard.server.common.data.device.profile.lwm2m.TelemetryObserveStrategy.SINGLE; |
|
|
|
import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH; |
|
|
|
import static org.thingsboard.server.common.data.util.CollectionsUtil.diffSets; |
|
|
|
import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaUpdateService.FW_3_VER_ID; |
|
|
|
@ -475,7 +480,7 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl |
|
|
|
Set<String> supportedObjects = clientContext.getSupportedIdVerInClient(lwM2MClient); |
|
|
|
if (supportedObjects != null && supportedObjects.size() > 0) { |
|
|
|
this.sendReadRequests(lwM2MClient, profile, supportedObjects); |
|
|
|
this.sendObserveRequests(lwM2MClient, profile, supportedObjects); |
|
|
|
this.sendInitObserveRequests(lwM2MClient, profile, supportedObjects); |
|
|
|
this.sendWriteAttributeRequests(lwM2MClient, profile, supportedObjects); |
|
|
|
// Removed. Used only for debug.
|
|
|
|
// this.sendDiscoverRequests(lwM2MClient, profile, supportedObjects);
|
|
|
|
@ -501,16 +506,27 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void sendObserveRequests(LwM2mClient lwM2MClient, Lwm2mDeviceProfileTransportConfiguration profile, Set<String> supportedObjects) { |
|
|
|
private void sendInitObserveRequests(LwM2mClient lwM2MClient, Lwm2mDeviceProfileTransportConfiguration profile, Set<String> supportedObjects) { |
|
|
|
try { |
|
|
|
Set<String> targetIds = profile.getObserveAttr().getObserve(); |
|
|
|
targetIds = targetIds.stream().filter(target -> isSupportedTargetId(supportedObjects, target)).collect(Collectors.toSet()); |
|
|
|
|
|
|
|
CountDownLatch latch = new CountDownLatch(targetIds.size()); |
|
|
|
targetIds.forEach(targetId -> sendObserveRequest(lwM2MClient, targetId, |
|
|
|
new TbLwM2MLatchCallback<>(latch, new TbLwM2MObserveCallback(this, logService, lwM2MClient, targetId)))); |
|
|
|
|
|
|
|
latch.await(config.getTimeout(), TimeUnit.MILLISECONDS); |
|
|
|
if (!targetIds.isEmpty()) { |
|
|
|
TelemetryObserveStrategy observeStrategy = profile.getObserveAttr().getObserveStrategy(); |
|
|
|
if (SINGLE.equals(observeStrategy)) { |
|
|
|
CountDownLatch latch = new CountDownLatch(targetIds.size()); |
|
|
|
targetIds.forEach(targetId -> sendObserveRequest(lwM2MClient, targetId, |
|
|
|
new TbLwM2MLatchCallback<>(latch, new TbLwM2MObserveCallback(this, logService, lwM2MClient, targetId)))); |
|
|
|
latch.await(config.getTimeout(), TimeUnit.MILLISECONDS); |
|
|
|
} else if (COMPOSITE_ALL.equals(observeStrategy)) { |
|
|
|
String[] versionedIds = targetIds.toArray(new String[0]); |
|
|
|
sendObserveCompositeRequest(lwM2MClient, versionedIds); |
|
|
|
} else if (COMPOSITE_BY_OBJECT.equals(observeStrategy)) { |
|
|
|
Map<Integer, String[]> versionedObjectIds = groupByObjectIdVersionedIds(targetIds); |
|
|
|
CountDownLatch latch = new CountDownLatch(versionedObjectIds.size()); |
|
|
|
versionedObjectIds.forEach((k, v)-> sendObserveCompositeRequest(lwM2MClient, v)); |
|
|
|
latch.await(config.getTimeout(), TimeUnit.MILLISECONDS); |
|
|
|
} |
|
|
|
} |
|
|
|
} catch (InterruptedException e) { |
|
|
|
log.error("[{}] Failed to await Observe requests!", lwM2MClient.getEndpoint(), e); |
|
|
|
} catch (Exception e) { |
|
|
|
@ -547,6 +563,12 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl |
|
|
|
TbLwM2MObserveRequest request = TbLwM2MObserveRequest.builder().versionedId(versionedId).timeout(clientContext.getRequestTimeout(lwM2MClient)).build(); |
|
|
|
defaultLwM2MDownlinkMsgHandler.sendObserveRequest(lwM2MClient, request, callback); |
|
|
|
} |
|
|
|
private void sendObserveCompositeRequest(LwM2mClient lwM2MClient, String[] versionedIds) { |
|
|
|
|
|
|
|
TbLwM2MObserveCompositeRequest request = TbLwM2MObserveCompositeRequest.builder().versionedIds(versionedIds).timeout(clientContext.getRequestTimeout(lwM2MClient)).build(); |
|
|
|
var mainCallback = new TbLwM2MObserveCompositeCallback(this, logService, lwM2MClient, versionedIds); |
|
|
|
defaultLwM2MDownlinkMsgHandler.sendObserveCompositeRequest(lwM2MClient, request, mainCallback); |
|
|
|
} |
|
|
|
|
|
|
|
private void sendWriteAttributesRequest(LwM2mClient lwM2MClient, String targetId, ObjectAttributes params) { |
|
|
|
TbLwM2MWriteAttributesRequest request = TbLwM2MWriteAttributesRequest.builder().versionedId(targetId).attributes(params).timeout(clientContext.getRequestTimeout(lwM2MClient)).build(); |
|
|
|
@ -1024,4 +1046,14 @@ public class DefaultLwM2mUplinkMsgHandler extends LwM2MExecutorAwareService impl |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private Map<Integer, String[]> groupByObjectIdVersionedIds(Set<String> targetIds){ |
|
|
|
return targetIds.stream() |
|
|
|
.collect(Collectors.groupingBy( |
|
|
|
id -> new LwM2mPath(fromVersionedIdToObjectId(id)).getObjectId(), |
|
|
|
Collectors.collectingAndThen( |
|
|
|
Collectors.toList(), |
|
|
|
list -> list.toArray(new String[0]) |
|
|
|
) |
|
|
|
)); |
|
|
|
} |
|
|
|
} |
|
|
|
|