diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/lwm2m/OtherConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/device/data/lwm2m/OtherConfiguration.java index e45a3d0e31..81d9d56a8c 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/device/data/lwm2m/OtherConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/device/data/lwm2m/OtherConfiguration.java @@ -27,5 +27,6 @@ public class OtherConfiguration { private PowerMode powerMode; private String fwUpdateResource; private String swUpdateResource; + private boolean composite; } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/config/LwM2mVersion.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/config/LwM2mVersion.java new file mode 100644 index 0000000000..079dc296d9 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/config/LwM2mVersion.java @@ -0,0 +1,66 @@ +/** + * 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.config; + +import lombok.Getter; +import org.eclipse.leshan.core.LwM2m.Version; +import org.eclipse.leshan.core.request.ContentFormat; + +public enum LwM2mVersion { + VERSION_1_0(0, Version.V1_0, ContentFormat.TLV), + VERSION_1_1(1, Version.V1_1, ContentFormat.TEXT); + + @Getter + private final int code; + @Getter + private final Version version; + @Getter + private final ContentFormat contentFormat; + + LwM2mVersion(int code, Version version, ContentFormat contentFormat) { + this.code = code; + this.version = version; + this.contentFormat = contentFormat; + } + + public static LwM2mVersion fromVersion(Version version) { + for (LwM2mVersion to : LwM2mVersion.values()) { + if (to.version.equals(version)) { + return to; + } + } + throw new IllegalArgumentException(String.format("Unsupported typeLwM2mVersion type : %s", version)); + } + + public static LwM2mVersion fromVersionStr(String versionStr) { + for (LwM2mVersion to : LwM2mVersion.values()) { + if (to.version.toString().equals(versionStr)) { + return to; + } + } + throw new IllegalArgumentException(String.format("Unsupported contentFormatLwM2mVersion version : %s", versionStr)); + } + + public static LwM2mVersion fromCode(int code) { + for (LwM2mVersion to : LwM2mVersion.values()) { + if (to.code == code) { + return to; + } + } + throw new IllegalArgumentException(String.format("Unsupported codeLwM2mVersion code : %d", code)); + } +} + diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mOperationType.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mOperationType.java index f9b3f93854..6cd204299a 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mOperationType.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mOperationType.java @@ -23,37 +23,41 @@ import lombok.Getter; public enum LwM2mOperationType { READ(0, "Read", true), - DISCOVER(1, "Discover", true), - DISCOVER_ALL(2, "DiscoverAll", false), - OBSERVE_READ_ALL(3, "ObserveReadAll", false), + READ_COMPOSITE(1, "ReadComposite", false, true), + DISCOVER(2, "Discover", true), + DISCOVER_ALL(3, "DiscoverAll", false), + OBSERVE_READ_ALL(4, "ObserveReadAll", false), - OBSERVE(4, "Observe", true), - OBSERVE_CANCEL(5, "ObserveCancel", true), - OBSERVE_CANCEL_ALL(6, "ObserveCancelAll", false), - EXECUTE(7, "Execute", true), + OBSERVE(5, "Observe", true), + OBSERVE_COMPOSITE(6, "ObserveComposite", false, true), + OBSERVE_CANCEL(7, "ObserveCancel", true), + OBSERVE_COMPOSITE_CANCEL(8, "ObserveCompositeCancel", false, true), + OBSERVE_CANCEL_ALL(9, "ObserveCancelAll", false), + EXECUTE(10, "Execute", true), /** * Replaces the Object Instance or the Resource(s) with the new value provided in the “Write” operation. (see * section 5.3.3 of the LW M2M spec). * if all resources are to be replaced */ - WRITE_REPLACE(8, "WriteReplace", true), + WRITE_REPLACE(11, "WriteReplace", true), /** * Adds or updates Resources provided in the new value and leaves other existing Resources unchanged. (see section * 5.3.3 of the LW M2M spec). * if this is a partial update request */ - WRITE_UPDATE(9, "WriteUpdate", true), - WRITE_ATTRIBUTES(10, "WriteAttributes", true), - DELETE(11, "Delete", true), + WRITE_UPDATE(12, "WriteUpdate", true), + WRITE_COMPOSITE(14, "WriteComposite", false, true), + WRITE_ATTRIBUTES(15, "WriteAttributes", true), + DELETE(16, "Delete", true), // only for RPC - FW_UPDATE(12, "FirmwareUpdate", false); + FW_UPDATE(17, "FirmwareUpdate", false); -// FW_READ_INFO(12, "FirmwareReadInfo"), -// SW_READ_INFO(15, "SoftwareReadInfo"), -// SW_UPDATE(16, "SoftwareUpdate"), -// SW_UNINSTALL(18, "SoftwareUninstall"); +// FW_READ_INFO(18, "FirmwareReadInfo"), +// SW_READ_INFO(19, "SoftwareReadInfo"), +// SW_UPDATE(20, "SoftwareUpdate"), +// SW_UNINSTALL(21, "SoftwareUninstall"); @Getter private final int code; @@ -62,10 +66,21 @@ public enum LwM2mOperationType { @Getter private final boolean hasObjectId; + @Getter + private final boolean composite; + LwM2mOperationType(int code, String type, boolean hasObjectId) { + this(code, type, hasObjectId, false); + } + + LwM2mOperationType(int code, String type, boolean hasObjectId, boolean composite) { this.code = code; this.type = type; this.hasObjectId = hasObjectId; + this.composite = composite; + if(hasObjectId && composite){ + throw new IllegalArgumentException("Can't set both Composite and hasObjectId for the same operation!"); + } } public static LwM2mOperationType fromType(String type) { diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java index 59df37de48..923c9d7bad 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mTransportUtil.java @@ -44,6 +44,7 @@ import org.thingsboard.server.common.data.device.data.lwm2m.BootstrapConfigurati import org.thingsboard.server.common.data.device.profile.DeviceProfileTransportConfiguration; import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration; import org.thingsboard.server.common.transport.TransportServiceCallback; +import org.thingsboard.server.transport.lwm2m.config.LwM2mVersion; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import org.thingsboard.server.transport.lwm2m.server.client.ResourceValue; import org.thingsboard.server.transport.lwm2m.server.ota.firmware.FirmwareUpdateResult; @@ -85,7 +86,7 @@ import static org.thingsboard.server.transport.lwm2m.server.ota.DefaultLwM2MOtaU @Slf4j public class LwM2mTransportUtil { - public static final String LWM2M_VERSION_DEFAULT = "1.0"; + public static final String LWM2M_OBJECT_VERSION_DEFAULT = "1.0"; public static final String LOG_LWM2M_TELEMETRY = "transportLog"; public static final String LOG_LWM2M_INFO = "info"; @@ -234,7 +235,7 @@ public class LwM2mTransportUtil { public static String convertObjectIdToVersionedId(String path, Registration registration) { String ver = registration.getSupportedObject().get(new LwM2mPath(path).getObjectId()); - ver = ver != null ? ver : LWM2M_VERSION_DEFAULT; + ver = ver != null ? ver : LwM2mVersion.VERSION_1_0.getVersion().toString(); try { String[] keyArray = path.split(LWM2M_SEPARATOR_PATH); if (keyArray.length > 1) { @@ -381,4 +382,5 @@ public class LwM2mTransportUtil { } return lwm2mResourceValue; } + } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java index 15e023404c..8cf6f20f36 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java @@ -37,6 +37,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.transport.auth.ValidateDeviceCredentialsResponse; import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto; import org.thingsboard.server.gen.transport.TransportProtos.TsKvProto; +import org.thingsboard.server.transport.lwm2m.config.LwM2mVersion; import java.io.IOException; import java.io.ObjectInputStream; @@ -52,7 +53,7 @@ import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_PATH; -import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LWM2M_VERSION_DEFAULT; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LWM2M_OBJECT_VERSION_DEFAULT; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertObjectIdToVersionedId; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.equalsResourceTypeGetSimpleName; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId; @@ -226,6 +227,10 @@ public class LwM2mClient implements Serializable { .getResourceModel(pathIds.getObjectId(), pathIds.getResourceId()) : null; } + public boolean isResourceMultiInstances(String pathIdVer, LwM2mModelProvider modelProvider) { + return getResourceModel(pathIdVer, modelProvider).multiple; + } + public ObjectModel getObjectModel(String pathIdVer, LwM2mModelProvider modelProvider) { LwM2mPath pathIds = new LwM2mPath(fromVersionedIdToObjectId(pathIdVer)); String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId()); @@ -234,11 +239,11 @@ public class LwM2mClient implements Serializable { .getObjectModel(pathIds.getObjectId()) : null; } - public String objectToString(LwM2mObject lwM2mObject, LwM2mValueConverter converter, String pathIdVer) { + public String objectToString(LwM2mObject lwM2mObject) { StringBuilder builder = new StringBuilder(); builder.append("LwM2mObject [id=").append(lwM2mObject.getId()).append(", instances={"); lwM2mObject.getInstances().forEach((instId, inst) -> { - builder.append(instId).append("=").append(this.instanceToString(inst, converter, pathIdVer)).append(", "); + builder.append(instId).append("=").append(this.instanceToString(inst)).append(", "); }); int startInd = builder.lastIndexOf(", "); if (startInd > 0) { @@ -248,11 +253,11 @@ public class LwM2mClient implements Serializable { return builder.toString(); } - public String instanceToString(LwM2mObjectInstance objectInstance, LwM2mValueConverter converter, String pathIdVer) { + public String instanceToString(LwM2mObjectInstance objectInstance) { StringBuilder builder = new StringBuilder(); builder.append("LwM2mObjectInstance [id=").append(objectInstance.getId()).append(", resources={"); objectInstance.getResources().forEach((resId, res) -> { - builder.append(resId).append("=").append(this.resourceToString(res, converter, pathIdVer)).append(", "); + builder.append(resId).append("=").append(this.resourceToString(res)).append(", "); }); int startInd = builder.lastIndexOf(", "); if (startInd > 0) { @@ -262,8 +267,8 @@ public class LwM2mClient implements Serializable { return builder.toString(); } - public String resourceToString(LwM2mResource lwM2mResource, LwM2mValueConverter converter, String pathIdVer) { - return lwM2mResource.getValue().toString(); + public String resourceToString(LwM2mResource lwM2mResource) { + return lwM2mResource.isMultiInstances() ? lwM2mResource.getInstances().toString() : lwM2mResource.getValue().toString(); } public Collection getNewResourceForInstance(String pathRezIdVer, Object params, LwM2mModelProvider modelProvider, @@ -299,11 +304,14 @@ public class LwM2mClient implements Serializable { return resources; } - public boolean isValidObjectVersion(String path) { + public void isValidObjectVersion(String path) { LwM2mPath pathIds = new LwM2mPath(fromVersionedIdToObjectId(path)); String verSupportedObject = registration.getSupportedObject().get(pathIds.getObjectId()); String verRez = getVerFromPathIdVerOrId(path); - return verRez == null ? LWM2M_VERSION_DEFAULT.equals(verSupportedObject) : verRez.equals(verSupportedObject); + if ((verRez != null && !verRez.equals(verSupportedObject)) || + (verRez == null && !LWM2M_OBJECT_VERSION_DEFAULT.equals(verSupportedObject))) { + throw new IllegalArgumentException(String.format("Specified resource id %s is not valid version! Must be version: %s", path, verSupportedObject)); + } } /** @@ -348,10 +356,8 @@ public class LwM2mClient implements Serializable { public ContentFormat getDefaultContentFormat() { if (registration == null) { return ContentFormat.DEFAULT; - } else if (registration.getLwM2mVersion().equals("1.0")) { - return ContentFormat.TLV; } else { - return ContentFormat.TEXT; + return LwM2mVersion.fromVersionStr(registration.getLwM2mVersion()).getContentFormat(); } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java index 238f948939..1599d73104 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java @@ -26,24 +26,31 @@ import org.eclipse.leshan.core.node.LwM2mResource; import org.eclipse.leshan.core.node.ObjectLink; import org.eclipse.leshan.core.node.codec.CodecException; import org.eclipse.leshan.core.observation.Observation; +import org.eclipse.leshan.core.request.CompositeDownlinkRequest; import org.eclipse.leshan.core.request.ContentFormat; import org.eclipse.leshan.core.request.DeleteRequest; import org.eclipse.leshan.core.request.DiscoverRequest; +import org.eclipse.leshan.core.request.DownlinkRequest; import org.eclipse.leshan.core.request.ExecuteRequest; import org.eclipse.leshan.core.request.ObserveRequest; +import org.eclipse.leshan.core.request.ReadCompositeRequest; import org.eclipse.leshan.core.request.ReadRequest; import org.eclipse.leshan.core.request.SimpleDownlinkRequest; import org.eclipse.leshan.core.request.WriteAttributesRequest; +import org.eclipse.leshan.core.request.WriteCompositeRequest; import org.eclipse.leshan.core.request.WriteRequest; import org.eclipse.leshan.core.response.DeleteResponse; import org.eclipse.leshan.core.response.DiscoverResponse; import org.eclipse.leshan.core.response.ExecuteResponse; import org.eclipse.leshan.core.response.LwM2mResponse; import org.eclipse.leshan.core.response.ObserveResponse; +import org.eclipse.leshan.core.response.ReadCompositeResponse; import org.eclipse.leshan.core.response.ReadResponse; import org.eclipse.leshan.core.response.WriteAttributesResponse; +import org.eclipse.leshan.core.response.WriteCompositeResponse; import org.eclipse.leshan.core.response.WriteResponse; import org.eclipse.leshan.core.util.Hex; +import org.eclipse.leshan.server.model.LwM2mModelProvider; import org.eclipse.leshan.server.registration.Registration; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; @@ -53,7 +60,10 @@ import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig; import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import org.thingsboard.server.transport.lwm2m.server.common.LwM2MExecutorAwareService; +import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MReadCompositeRequest; +import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MWriteResponseCompositeCallback; import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; +import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; import javax.annotation.PostConstruct; @@ -63,6 +73,7 @@ import java.util.Collection; import java.util.Date; import java.util.LinkedList; import java.util.List; +import java.util.Map; import java.util.Set; import java.util.function.Function; import java.util.function.Predicate; @@ -73,6 +84,7 @@ import static org.eclipse.leshan.core.attributes.Attribute.LESSER_THAN; import static org.eclipse.leshan.core.attributes.Attribute.MAXIMUM_PERIOD; import static org.eclipse.leshan.core.attributes.Attribute.MINIMUM_PERIOD; import static org.eclipse.leshan.core.attributes.Attribute.STEP; +import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId; @Slf4j @Service @@ -110,8 +122,31 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im @Override public void sendReadRequest(LwM2mClient client, TbLwM2MReadRequest request, DownlinkRequestCallback callback) { validateVersionedId(client, request); - ReadRequest downlink = new ReadRequest(getContentFormat(client, request), request.getObjectId()); - sendRequest(client, downlink, request.getTimeout(), callback); + ReadRequest downlink = new ReadRequest(getRequestContentFormat(client, request, this.config.getModelProvider()), request.getObjectId()); + sendSimpleRequest(client, downlink, request.getTimeout(), callback); + } + + @Override + public void sendReadCompositeRequest(LwM2mClient client, TbLwM2MReadCompositeRequest request, DownlinkRequestCallback callback) { + validateVersionedIds(client, request); + ContentFormat requestContentFormat = ContentFormat.SENML_JSON; + ContentFormat responseContentFormat = ContentFormat.SENML_JSON; + + ReadCompositeRequest downlink = new ReadCompositeRequest(requestContentFormat, responseContentFormat, request.getObjectIds()); + sendCompositeRequest(client, downlink, this.config.getTimeout(), callback); + } + + @Override + public void sendWriteCompositeRequest(LwM2mClient client, Map nodes, DefaultLwM2MUplinkMsgHandler handler) { +// ResourceModel resourceModelWrite = client.getResourceModel(request.getVersionedId(), this.config.getModelProvider()); + TbLwM2MWriteResponseCompositeCallback callback = new TbLwM2MWriteResponseCompositeCallback(handler, logService, client, null); + ContentFormat contentFormat = ContentFormat.SENML_JSON; + try { + WriteCompositeRequest downlink = new WriteCompositeRequest(contentFormat, nodes); + sendWriteCompositeRequest(client, downlink, this.config.getTimeout(), callback); + } catch (Exception e) { + callback.onError(JacksonUtil.toString(nodes), e); + } } @Override @@ -121,7 +156,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im Set observations = context.getServer().getObservationService().getObservations(client.getRegistration()); if (observations.stream().noneMatch(observation -> observation.getPath().equals(resultIds))) { ObserveRequest downlink; - ContentFormat contentFormat = getContentFormat(client, request); + ContentFormat contentFormat = getRequestContentFormat(client, request, this.config.getModelProvider()); if (resultIds.isResource()) { downlink = new ObserveRequest(contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resultIds.getResourceId()); } else if (resultIds.isObjectInstance()) { @@ -130,7 +165,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im downlink = new ObserveRequest(contentFormat, resultIds.getObjectId()); } log.info("[{}] Send observation: {}.", client.getEndpoint(), request.getVersionedId()); - sendRequest(client, downlink, request.getTimeout(), callback); + sendSimpleRequest(client, downlink, request.getTimeout(), callback); } else { callback.onValidationError(resultIds.toString(), "Observation is already registered!"); } @@ -158,13 +193,13 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im } else { downlink = new ExecuteRequest(request.getVersionedId()); } - sendRequest(client, downlink, request.getTimeout(), callback); + sendSimpleRequest(client, downlink, request.getTimeout(), callback); } } @Override public void sendDeleteRequest(LwM2mClient client, TbLwM2MDeleteRequest request, DownlinkRequestCallback callback) { - sendRequest(client, new DeleteRequest(request.getObjectId()), request.getTimeout(), callback); + sendSimpleRequest(client, new DeleteRequest(request.getObjectId()), request.getTimeout(), callback); } @Override @@ -182,7 +217,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im @Override public void sendDiscoverRequest(LwM2mClient client, TbLwM2MDiscoverRequest request, DownlinkRequestCallback callback) { validateVersionedId(client, request); - sendRequest(client, new DiscoverRequest(request.getObjectId()), request.getTimeout(), callback); + sendSimpleRequest(client, new DiscoverRequest(request.getObjectId()), request.getTimeout(), callback); } @Override @@ -202,7 +237,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im addAttribute(attributes, LESSER_THAN, params.getLt()); addAttribute(attributes, STEP, params.getSt()); AttributeSet attributeSet = new AttributeSet(attributes); - sendRequest(client, new WriteAttributesRequest(request.getObjectId(), attributeSet), request.getTimeout(), callback); + sendSimpleRequest(client, new WriteAttributesRequest(request.getObjectId(), attributeSet), request.getTimeout(), callback); } @Override @@ -214,7 +249,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im LwM2mPath path = new LwM2mPath(request.getObjectId()); WriteRequest downlink = this.getWriteRequestSingleResource(resourceModelWrite.type, contentFormat, path.getObjectId(), path.getObjectInstanceId(), path.getResourceId(), request.getValue()); - sendRequest(client, downlink, request.getTimeout(), callback); + sendSimpleRequest(client, downlink, request.getTimeout(), callback); } catch (Exception e) { callback.onError(JacksonUtil.toString(request), e); } @@ -236,7 +271,7 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im ContentFormat contentFormat = request.getObjectContentFormat() != null ? request.getObjectContentFormat() : convertResourceModelTypeToContentFormat(client, resourceModelWrite.type); WriteRequest downlink = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resources); - sendRequest(client, downlink, request.getTimeout(), callback); + sendSimpleRequest(client, downlink, request.getTimeout(), callback); } else if (resultIds.isObjectInstance()) { /* * params = "{\"id\":0,\"resources\":[{\"id\":14,\"value\":\"+5\"},{\"id\":15,\"value\":\"+9\"}]}" @@ -245,9 +280,9 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im */ Collection resources = client.getNewResourcesForInstance(request.getVersionedId(), request.getValue(), this.config.getModelProvider(), this.converter); if (resources.size() > 0) { - ContentFormat contentFormat = request.getObjectContentFormat() != null ? request.getObjectContentFormat() : client.getDefaultContentFormat(); + ContentFormat contentFormat = request.getObjectContentFormat() != null ? request.getObjectContentFormat() : ContentFormat.DEFAULT; WriteRequest downlink = new WriteRequest(WriteRequest.Mode.UPDATE, contentFormat, resultIds.getObjectId(), resultIds.getObjectInstanceId(), resources); - sendRequest(client, downlink, request.getTimeout(), callback); + sendSimpleRequest(client, downlink, request.getTimeout(), callback); } else { callback.onValidationError(JacksonUtil.toString(request), "No resources to update!"); } @@ -256,10 +291,42 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im } } - private , T extends LwM2mResponse> void sendRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback callback) { + + private , T extends LwM2mResponse> void sendSimpleRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback callback) { + sendRequest(client, request, timeoutInMs, callback, r -> request.getPath().toString()); + } + + private , T extends LwM2mResponse> void sendCompositeRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback callback) { + sendRequest(client, request, timeoutInMs, callback, r -> request.getPaths().toString()); + } + + private , T extends LwM2mResponse> void sendRequest(LwM2mClient client, R request, long timeoutInMs, DownlinkRequestCallback callback, Function pathToStringFunction) { + Registration registration = client.getRegistration(); + try { + logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), pathToStringFunction.apply(request))); + context.getServer().send(registration, request, timeoutInMs, response -> { + executor.submit(() -> { + try { + callback.onSuccess(request, response); + } catch (Exception e) { + log.error("[{}] failed to process successful response [{}] ", registration.getEndpoint(), response, e); + } + }); + }, e -> { + executor.submit(() -> { + callback.onError(JacksonUtil.toString(request), e); + }); + }); + } catch (Exception e) { + callback.onError(JacksonUtil.toString(request), e); + } + } + + private , T extends LwM2mResponse> void sendWriteCompositeRequest(LwM2mClient client, WriteCompositeRequest request, long timeoutInMs, + DownlinkRequestCallback callback) { Registration registration = client.getRegistration(); try { - logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), request.getPath())); + logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), request.getPaths())); context.getServer().send(registration, request, timeoutInMs, response -> { executor.submit(() -> { try { @@ -278,6 +345,52 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im } } +// private , T extends LwM2mResponse> void sendReadRequestComposite(LwM2mClient client, ReadCompositeRequest request, long timeoutInMs, +// DownlinkRequestCallback callback) { +// Registration registration = client.getRegistration(); +// try { +// logService.log(client, String.format("[%s][%s] Sending request: %s to %s", registration.getId(), registration.getSocketAddress(), request.getClass().getSimpleName(), request.getPaths())); +// context.getServer().send(registration, request, timeoutInMs, response -> { +// executor.submit(() -> { +// try { +// /** +// * [{"bn":"/3/0/","n":"0","vs":"Thingsboard Test Device"}, +// * {"n":"1","vs":"Model 500"}, +// * {"n":"2","vs":"TH-500-000-0001"}, +// * {"n":"3","vs":"TestThingsboard@TestMore1024_2.04"}, +// * {"n":"6","v":1},{"n":"7","v":56}, +// * {"n":"8","v":42},{"n":"9","v":16}, +// * {"n":"10","v":127619},{"n":"13","v":1624520988}, +// * {"n":"14","vs":"+03"},{"n":"15","vs":"Europe/Kiev"}, +// * {"n":"16","vs":"U"},{"n":"17","vs":"smart meters"}, +// * {"n":"18","vs":"1.01"},{"n":"19","vs":"1.02"}, +// * {"n":"20","v":3},{"n":"21","v":256000}, +// * {"bn":"/5/0/","n":"1","vs":""}, +// * {"n":"3","v":0},{"n":"5","v":0}, +// * {"n":"6","vs":""},{"n":"7","vs":""}, +// * {"n":"8/0","v":0},{"n":"8/1","v":1}, +// * {"n":"9","v":2}, +// * {"bn":"/1/0/","n":"0","v":123}, +// * {"n":"1","v":300}, +// * {"n":"6","vb":false}, +// * {"n":"22","vs":"U"}, +// * {"n":"7","vs":"U"}] +// */ +// callback.onSuccess(request, response); +// } catch (Exception e) { +// log.error("[{}] failed to process successful response [{}] ", registration.getEndpoint(), response, e); +// } +// }); +// }, e -> { +// executor.submit(() -> { +// callback.onError(JacksonUtil.toString(request), e); +// }); +// }); +// } catch (Exception e) { +// callback.onError(JacksonUtil.toString(request), e); +// } +// } + private WriteRequest getWriteRequestSingleResource(ResourceModel.Type type, ContentFormat contentFormat, int objectId, int instanceId, int resourceId, Object value) { switch (type) { case STRING: // String @@ -308,14 +421,23 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im } private void validateVersionedId(LwM2mClient client, HasVersionedId request) { - if (!client.isValidObjectVersion(request.getVersionedId())) { - throw new IllegalArgumentException("Specified resource id is not configured in the device profile!"); - } + client.isValidObjectVersion(request.getVersionedId()); if (request.getObjectId() == null) { throw new IllegalArgumentException("Specified object id is null!"); } } + private void validateVersionedIds(LwM2mClient client, HasVersionedIds request) { + for (String versionedId : request.getVersionedIds()) { + client.isValidObjectVersion(versionedId); + } + for (String objectId : request.getObjectIds()) { + if (objectId == null) { + throw new IllegalArgumentException("Specified object id is null!"); + } + } + } + private static void addAttribute(List attributes, String attributeName, T value) { addAttribute(attributes, attributeName, value, null, null); } @@ -347,7 +469,23 @@ public class DefaultLwM2mDownlinkMsgHandler extends LwM2MExecutorAwareService im throw new CodecException("Invalid ResourceModel_Type for %s ContentFormat.", type); } - private static ContentFormat getContentFormat(LwM2mClient client, HasContentFormat request) { - return request.getContentFormat() != null ? request.getContentFormat() : client.getDefaultContentFormat(); + private static ContentFormat getRequestContentFormat(LwM2mClient client, HasContentFormat request, LwM2mModelProvider modelProvider) { + if (request.getRequestContentFormat() != null) { + return request.getRequestContentFormat(); + } else { + String versionedId = null; + if (request instanceof TbLwM2MReadRequest) { + versionedId = ((TbLwM2MReadRequest) request).getVersionedId(); + } else if (request instanceof TbLwM2MObserveRequest) { + versionedId = ((TbLwM2MObserveRequest) request).getVersionedId(); + } + String id = fromVersionedIdToObjectId(versionedId); + if (id != null && new LwM2mPath(id).isResource() && !client.isResourceMultiInstances(versionedId, modelProvider)) { + return client.getDefaultContentFormat(); + } + else { + return ContentFormat.DEFAULT; + } + } } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/HasContentFormat.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/HasContentFormat.java index 6d2ecc98ed..f2618224e7 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/HasContentFormat.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/HasContentFormat.java @@ -19,5 +19,9 @@ import org.eclipse.leshan.core.request.ContentFormat; public interface HasContentFormat { - ContentFormat getContentFormat(); + ContentFormat getRequestContentFormat(); + + default ContentFormat getResponseContentFormat() { + return null; + } } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/HasVersionedIds.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/HasVersionedIds.java new file mode 100644 index 0000000000..d0b920ca1c --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/HasVersionedIds.java @@ -0,0 +1,35 @@ +/** + * 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.downlink; + +import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil; + +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + +public interface HasVersionedIds { + + String[] getVersionedIds(); + + default String[] getObjectIds() { + Set objectIds = ConcurrentHashMap.newKeySet(); + for (String versionedId : getVersionedIds()) { + objectIds.add(LwM2mTransportUtil.fromVersionedIdToObjectId(versionedId)); + } + return (String[]) objectIds.toArray(String[]::new); + } + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/LwM2mDownlinkMsgHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/LwM2mDownlinkMsgHandler.java index adb2e32293..04c63aec17 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/LwM2mDownlinkMsgHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/LwM2mDownlinkMsgHandler.java @@ -20,6 +20,7 @@ import org.eclipse.leshan.core.request.DeleteRequest; import org.eclipse.leshan.core.request.DiscoverRequest; import org.eclipse.leshan.core.request.ExecuteRequest; import org.eclipse.leshan.core.request.ObserveRequest; +import org.eclipse.leshan.core.request.ReadCompositeRequest; import org.eclipse.leshan.core.request.ReadRequest; import org.eclipse.leshan.core.request.WriteAttributesRequest; import org.eclipse.leshan.core.request.WriteRequest; @@ -27,18 +28,24 @@ import org.eclipse.leshan.core.response.DeleteResponse; import org.eclipse.leshan.core.response.DiscoverResponse; import org.eclipse.leshan.core.response.ExecuteResponse; import org.eclipse.leshan.core.response.ObserveResponse; +import org.eclipse.leshan.core.response.ReadCompositeResponse; import org.eclipse.leshan.core.response.ReadResponse; import org.eclipse.leshan.core.response.WriteAttributesResponse; import org.eclipse.leshan.core.response.WriteResponse; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; +import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MReadCompositeRequest; +import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler; import java.util.List; +import java.util.Map; import java.util.Set; public interface LwM2mDownlinkMsgHandler { void sendReadRequest(LwM2mClient client, TbLwM2MReadRequest request, DownlinkRequestCallback callback); + void sendReadCompositeRequest(LwM2mClient client, TbLwM2MReadCompositeRequest request, DownlinkRequestCallback callback); + void sendObserveRequest(LwM2mClient client, TbLwM2MObserveRequest request, DownlinkRequestCallback callback); void sendObserveAllRequest(LwM2mClient client, TbLwM2MObserveAllRequest request, DownlinkRequestCallback> callback); @@ -59,6 +66,8 @@ public interface LwM2mDownlinkMsgHandler { void sendWriteReplaceRequest(LwM2mClient client, TbLwM2MWriteReplaceRequest request, DownlinkRequestCallback callback); + void sendWriteCompositeRequest(LwM2mClient client, Map nodes, DefaultLwM2MUplinkMsgHandler handler); + void sendWriteUpdateRequest(LwM2mClient client, TbLwM2MWriteUpdateRequest request, DownlinkRequestCallback callback); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java index f3348aa2c9..0d13b4a1ff 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveRequest.java @@ -24,12 +24,12 @@ import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MObserveRequest extends AbstractTbLwM2MTargetedDownlinkRequest implements HasContentFormat { @Getter - private final ContentFormat contentFormat; + private final ContentFormat requestContentFormat; @Builder - private TbLwM2MObserveRequest(String versionedId, long timeout, ContentFormat contentFormat) { + private TbLwM2MObserveRequest(String versionedId, long timeout, ContentFormat requestContentFormat) { super(versionedId, timeout); - this.contentFormat = contentFormat; + this.requestContentFormat = requestContentFormat; } @Override diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java index a07e738465..6f388237c2 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadRequest.java @@ -24,12 +24,12 @@ import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; public class TbLwM2MReadRequest extends AbstractTbLwM2MTargetedDownlinkRequest implements HasContentFormat { @Getter - private final ContentFormat contentFormat; + private final ContentFormat requestContentFormat; @Builder - private TbLwM2MReadRequest(String versionedId, long timeout, ContentFormat contentFormat) { + private TbLwM2MReadRequest(String versionedId, long timeout, ContentFormat requestContentFormat) { super(versionedId, timeout); - this.contentFormat = contentFormat; + this.requestContentFormat = requestContentFormat; } @Override diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MTargetedCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MTargetedCallback.java index 373c882ec4..6b013ef9b5 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MTargetedCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MTargetedCallback.java @@ -18,7 +18,8 @@ package org.thingsboard.server.transport.lwm2m.server.downlink; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; -import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; + +import java.util.Arrays; import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_INFO; @@ -26,18 +27,26 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.L public abstract class TbLwM2MTargetedCallback extends AbstractTbLwM2MRequestCallback { protected final String versionedId; + protected final String[] versionedIds; public TbLwM2MTargetedCallback(LwM2MTelemetryLogService logService, LwM2mClient client, String versionedId) { super(logService, client); this.versionedId = versionedId; + this.versionedIds = null; + } + + public TbLwM2MTargetedCallback(LwM2MTelemetryLogService logService, LwM2mClient client, String[] versionedIds) { + super(logService, client); + this.versionedId = null; + this.versionedIds = versionedIds; } @Override public void onSuccess(R request, T response) { //TODO convert camelCase to "camel case" using .split("(? extends TbLwM2MTargetedCallback { @@ -32,4 +30,9 @@ public abstract class TbLwM2MUplinkTargetedCallback extends TbLwM2MTargete this.handler = handler; } + public TbLwM2MUplinkTargetedCallback(LwM2mUplinkMsgHandler handler, LwM2MTelemetryLogService logService, LwM2mClient client, String[] versionedIds) { + super(logService, client, versionedIds); + this.handler = handler; + } + } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/AbstractTbLwM2MTargetedDownlinkCompositeRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/AbstractTbLwM2MTargetedDownlinkCompositeRequest.java new file mode 100644 index 0000000000..1625e5aad0 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/AbstractTbLwM2MTargetedDownlinkCompositeRequest.java @@ -0,0 +1,34 @@ +/** + * 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.downlink.composite; + +import lombok.Getter; +import org.thingsboard.server.transport.lwm2m.server.downlink.HasVersionedIds; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MDownlinkRequest; + +public abstract class AbstractTbLwM2MTargetedDownlinkCompositeRequest implements TbLwM2MDownlinkRequest, HasVersionedIds { + + @Getter + private final String [] versionedIds; + @Getter + private final long timeout; + + public AbstractTbLwM2MTargetedDownlinkCompositeRequest(String [] versionedIds, long timeout) { + this.versionedIds = versionedIds; + this.timeout = timeout; + } + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MReadCompositeCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MReadCompositeCallback.java new file mode 100644 index 0000000000..82e0936b87 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MReadCompositeCallback.java @@ -0,0 +1,39 @@ +/** + * 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.downlink.composite; + +import lombok.extern.slf4j.Slf4j; +import org.eclipse.leshan.core.request.ReadCompositeRequest; +import org.eclipse.leshan.core.response.ReadCompositeResponse; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MUplinkTargetedCallback; +import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; + +@Slf4j +public class TbLwM2MReadCompositeCallback extends TbLwM2MUplinkTargetedCallback { + + public TbLwM2MReadCompositeCallback(LwM2mUplinkMsgHandler handler, LwM2MTelemetryLogService logService, LwM2mClient client, String[] versionedIds) { + super(handler, logService, client, versionedIds); + } + + @Override + public void onSuccess(ReadCompositeRequest request, ReadCompositeResponse response) { + super.onSuccess(request, response); + handler.onUpdateValueAfterReadCompositeResponse(client.getRegistration(), response); + } + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MReadCompositeRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MReadCompositeRequest.java new file mode 100644 index 0000000000..43b55f371f --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MReadCompositeRequest.java @@ -0,0 +1,44 @@ +/** + * 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.downlink.composite; + +import lombok.Builder; +import lombok.Getter; +import org.eclipse.leshan.core.request.ContentFormat; +import org.eclipse.leshan.core.response.ReadCompositeResponse; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; +import org.thingsboard.server.transport.lwm2m.server.downlink.HasContentFormat; + +public class TbLwM2MReadCompositeRequest extends AbstractTbLwM2MTargetedDownlinkCompositeRequest implements HasContentFormat { + + @Getter + private final ContentFormat requestContentFormat; + + @Getter + private final ContentFormat responseContentFormat; + + @Builder + private TbLwM2MReadCompositeRequest(String [] versionedIds, long timeout, ContentFormat requestContentFormat, ContentFormat responseContentFormat) { + super(versionedIds, timeout); + this.requestContentFormat = requestContentFormat; + this.responseContentFormat = responseContentFormat; + } + + @Override + public LwM2mOperationType getType() { + return LwM2mOperationType.READ_COMPOSITE; + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MWriteCompositeRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MWriteCompositeRequest.java new file mode 100644 index 0000000000..0e3ffaed28 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MWriteCompositeRequest.java @@ -0,0 +1,46 @@ +/** + * 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.downlink.composite; + +import lombok.Builder; +import lombok.Getter; +import org.eclipse.leshan.core.request.ContentFormat; +import org.eclipse.leshan.core.response.WriteCompositeResponse; +import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType; +import org.thingsboard.server.transport.lwm2m.server.downlink.AbstractTbLwM2MTargetedDownlinkRequest; + +public class TbLwM2MWriteCompositeRequest extends AbstractTbLwM2MTargetedDownlinkRequest { + + @Getter + private final ContentFormat contentFormat; + @Getter + private final Object value; + + @Builder + private TbLwM2MWriteCompositeRequest(String versionedId, long timeout, ContentFormat contentFormat, Object value) { + super(versionedId, timeout); + this.contentFormat = contentFormat; + this.value = value; + } + + @Override + public LwM2mOperationType getType() { + return LwM2mOperationType.WRITE_REPLACE; + } + + + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MWriteResponseCompositeCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MWriteResponseCompositeCallback.java new file mode 100644 index 0000000000..9dccc98185 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/composite/TbLwM2MWriteResponseCompositeCallback.java @@ -0,0 +1,37 @@ +/** + * 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.downlink.composite; + +import org.eclipse.leshan.core.request.WriteCompositeRequest; +import org.eclipse.leshan.core.response.WriteCompositeResponse; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; +import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MUplinkTargetedCallback; +import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; +import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; + +public class TbLwM2MWriteResponseCompositeCallback extends TbLwM2MUplinkTargetedCallback { + + public TbLwM2MWriteResponseCompositeCallback(LwM2mUplinkMsgHandler handler, LwM2MTelemetryLogService logService, LwM2mClient client, String targetId) { + super(handler, logService, client, targetId); + } + + @Override + public void onSuccess(WriteCompositeRequest request, WriteCompositeResponse response) { + super.onSuccess(request, response); + handler.onWriteCompositeResponseOk(client, request); + } + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java index 0432a5b742..f300e00675 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java @@ -50,7 +50,12 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttrib import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteReplaceRequest; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteResponseCallback; import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteUpdateRequest; +import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MReadCompositeCallback; +import org.thingsboard.server.transport.lwm2m.server.downlink.composite.TbLwM2MReadCompositeRequest; import org.thingsboard.server.transport.lwm2m.server.log.LwM2MTelemetryLogService; +import org.thingsboard.server.transport.lwm2m.server.rpc.composite.RpcReadCompositeRequest; +import org.thingsboard.server.transport.lwm2m.server.rpc.composite.RpcReadResponseCompositeCallback; +import org.thingsboard.server.transport.lwm2m.server.rpc.composite.RpcWriteCompositeRequest; import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler; import java.util.Map; @@ -125,6 +130,17 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { default: throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); } + } else if (operationType.isComposite()) { + switch (operationType) { + case READ_COMPOSITE: + sendReadCompositeRequest(client, rpcRequst); + break; + case WRITE_COMPOSITE: + sendWriteCompositeRequest(client, rpcRequst); + break; + default: + throw new IllegalArgumentException("Unsupported operation: " + operationType.name()); + } } else { switch (operationType) { case OBSERVE_CANCEL_ALL: @@ -151,14 +167,22 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { private void sendReadRequest(LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, String versionedId) { TbLwM2MReadRequest request = TbLwM2MReadRequest.builder().versionedId(versionedId).timeout(this.config.getTimeout()).build(); var mainCallback = new TbLwM2MReadCallback(uplinkHandler, logService, client, versionedId); - var rpcCallback = new RpcReadResponseCallback<>(transportService, client, requestMsg, versionedId, mainCallback); + var rpcCallback = new RpcReadResponseCallback<>(transportService, client, requestMsg, mainCallback); downlinkHandler.sendReadRequest(client, request, rpcCallback); } + private void sendReadCompositeRequest(LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg) { + String[] versionedIds = getIdsFromParameters(client, requestMsg); + TbLwM2MReadCompositeRequest request = TbLwM2MReadCompositeRequest.builder().versionedIds(versionedIds).timeout(this.config.getTimeout()).build(); + var mainCallback = new TbLwM2MReadCompositeCallback(uplinkHandler, logService, client, versionedIds); + var rpcCallback = new RpcReadResponseCompositeCallback(transportService, client, requestMsg, mainCallback); + downlinkHandler.sendReadCompositeRequest(client, request, rpcCallback); + } + private void sendObserveRequest(LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, String versionedId) { TbLwM2MObserveRequest request = TbLwM2MObserveRequest.builder().versionedId(versionedId).timeout(this.config.getTimeout()).build(); var mainCallback = new TbLwM2MObserveCallback(uplinkHandler, logService, client, versionedId); - var rpcCallback = new RpcReadResponseCallback<>(transportService, client, requestMsg, versionedId, mainCallback); + var rpcCallback = new RpcReadResponseCallback<>(transportService, client, requestMsg, mainCallback); downlinkHandler.sendObserveRequest(client, request, rpcCallback); } @@ -215,6 +239,16 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { downlinkHandler.sendWriteReplaceRequest(client, request, rpcCallback); } + private void sendWriteCompositeRequest(LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg) { + RpcWriteCompositeRequest nodes = JacksonUtil.fromString(requestMsg.getParams(), RpcWriteCompositeRequest.class); +// TbLwM2MWriteReplaceRequest request = TbLwM2MWriteReplaceRequest.builder().versionedId(versionedId) +// .value(requestBody.getValue()) +// .timeout(this.config.getTimeout()).build(); +// var mainCallback = new TbLwM2MWriteResponseCallback(uplinkHandler, logService, client, versionedId); +// var rpcCallback = new RpcEmptyResponseCallback<>(transportService, client, requestMsg, mainCallback); +// downlinkHandler.sendWriteReplaceRequest(client, request, rpcCallback); + } + private void sendCancelObserveRequest(LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, String versionedId) { TbLwM2MCancelObserveRequest downlink = TbLwM2MCancelObserveRequest.builder().versionedId(versionedId).timeout(this.config.getTimeout()).build(); var mainCallback = new TbLwM2MCancelObserveCallback(logService, client, versionedId); @@ -249,6 +283,24 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler { return targetId; } + private String[] getIdsFromParameters(LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg rpcRequst) { + RpcReadCompositeRequest requestParams = JacksonUtil.fromString(rpcRequst.getParams(), RpcReadCompositeRequest.class); + if (requestParams.getKeys() != null && requestParams.getKeys().length > 0) { + Set targetIds = ConcurrentHashMap.newKeySet(); + for (String key : requestParams.getKeys()) { + String targetId = clientContext.getObjectIdByKeyNameFromProfile(client, key); + if (targetId != null) { + targetIds.add(targetId); + } + } + return (String[]) targetIds.toArray(String[]::new); + } else if (requestParams.getIds() != null && requestParams.getIds().length > 0) { + return requestParams.getIds(); + } else { + throw new IllegalArgumentException("Can't find 'key' or 'id' in the requestParams parameters!"); + } + } + private void sendErrorRpcResponse(TransportProtos.SessionInfoProto sessionInfo, int requestId, String result, String error) { String payload = JacksonUtil.toString(JacksonUtil.newObjectNode().put("result", result).put("error", error)); TransportProtos.ToDeviceRpcResponseMsg msg = TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(payload).build(); diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/IdOrKeyRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/IdOrKeyRequest.java index 77d8cd3ea1..bef4fa37e8 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/IdOrKeyRequest.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/IdOrKeyRequest.java @@ -24,5 +24,4 @@ public class IdOrKeyRequest { private String key; private String id; - } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java index 115680b50f..c2e9a43ded 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcDownlinkRequestCallbackProxy.java @@ -16,13 +16,11 @@ package org.thingsboard.server.transport.lwm2m.server.rpc; import org.eclipse.leshan.core.ResponseCode; -import org.eclipse.leshan.core.node.codec.LwM2mValueConverter; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.transport.TransportService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; import org.thingsboard.server.transport.lwm2m.server.downlink.DownlinkRequestCallback; -import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl; public abstract class RpcDownlinkRequestCallbackProxy implements DownlinkRequestCallback { @@ -31,14 +29,12 @@ public abstract class RpcDownlinkRequestCallbackProxy implements DownlinkR private final DownlinkRequestCallback callback; protected final LwM2mClient client; - protected final LwM2mValueConverter converter; public RpcDownlinkRequestCallbackProxy(TransportService transportService, LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, DownlinkRequestCallback callback) { this.transportService = transportService; this.client = client; this.request = requestMsg; this.callback = callback; - this.converter = LwM2mValueConverterImpl.getInstance(); } @Override diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcReadResponseCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcReadResponseCallback.java index d3da6e6a3f..6f7eb4b2b7 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcReadResponseCallback.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/RpcReadResponseCallback.java @@ -29,22 +29,19 @@ import java.util.Optional; public class RpcReadResponseCallback, T extends ReadResponse> extends RpcLwM2MDownlinkCallback { - private final String versionedId; - - public RpcReadResponseCallback(TransportService transportService, LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, String versionedId, DownlinkRequestCallback callback) { + public RpcReadResponseCallback(TransportService transportService, LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, DownlinkRequestCallback callback) { super(transportService, client, requestMsg, callback); - this.versionedId = versionedId; } @Override protected Optional serializeSuccessfulResponse(T response) { Object value = null; if (response.getContent() instanceof LwM2mObject) { - value = client.objectToString((LwM2mObject) response.getContent(), this.converter, versionedId); + value = client.objectToString((LwM2mObject) response.getContent()); } else if (response.getContent() instanceof LwM2mObjectInstance) { - value = client.instanceToString((LwM2mObjectInstance) response.getContent(), this.converter, versionedId); + value = client.instanceToString((LwM2mObjectInstance) response.getContent()); } else if (response.getContent() instanceof LwM2mResource) { - value = client.resourceToString((LwM2mResource) response.getContent(), this.converter, versionedId); + value = client.resourceToString((LwM2mResource) response.getContent()); } return Optional.of(String.format("%s", value)); } diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/client/LwM2mSoftwareUpdate.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcReadCompositeRequest.java similarity index 70% rename from common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/client/LwM2mSoftwareUpdate.java rename to common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcReadCompositeRequest.java index ed8b316fed..ccd4767fee 100644 --- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/client/LwM2mSoftwareUpdate.java +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcReadCompositeRequest.java @@ -13,15 +13,16 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.thingsboard.server.transport.lwm2m.client; +package org.thingsboard.server.transport.lwm2m.server.rpc.composite; +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import lombok.Data; -import java.util.UUID; - @Data -public class LwM2mSoftwareUpdate { - private volatile String clientSwVersion; - private volatile String currentSwVersion; - private volatile UUID currentSwId; -} \ No newline at end of file +@JsonIgnoreProperties(ignoreUnknown = true) +public class RpcReadCompositeRequest { + + private String [] keys; + private String [] ids; + +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcReadResponseCompositeCallback.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcReadResponseCompositeCallback.java new file mode 100644 index 0000000000..8ce3a0a672 --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcReadResponseCompositeCallback.java @@ -0,0 +1,39 @@ +/** + * 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.rpc.composite; + +import org.eclipse.leshan.core.request.LwM2mRequest; +import org.eclipse.leshan.core.request.ReadCompositeRequest; +import org.eclipse.leshan.core.response.ReadCompositeResponse; +import org.thingsboard.server.common.transport.TransportService; +import org.thingsboard.server.gen.transport.TransportProtos; +import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient; +import org.thingsboard.server.transport.lwm2m.server.downlink.DownlinkRequestCallback; +import org.thingsboard.server.transport.lwm2m.server.rpc.RpcLwM2MDownlinkCallback; + +import java.util.Optional; + +public class RpcReadResponseCompositeCallback, T extends ReadCompositeResponse> extends RpcLwM2MDownlinkCallback { + + public RpcReadResponseCompositeCallback(TransportService transportService, LwM2mClient client, TransportProtos.ToDeviceRpcRequestMsg requestMsg, DownlinkRequestCallback callback) { + super(transportService, client, requestMsg, callback); + } + + @Override + protected Optional serializeSuccessfulResponse(T response) { + return Optional.of(String.format("%s", response.getContent().toString())); + } +} diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcWriteCompositeRequest.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcWriteCompositeRequest.java new file mode 100644 index 0000000000..ce13f58e9e --- /dev/null +++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/composite/RpcWriteCompositeRequest.java @@ -0,0 +1,29 @@ +/** + * 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.rpc.composite; + +import com.fasterxml.jackson.annotation.JsonIgnoreProperties; +import lombok.Data; + +import java.util.Map; + +@Data +@JsonIgnoreProperties(ignoreUnknown = true) +public class RpcWriteCompositeRequest { + + private Map nodes; + +} 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 9851612614..cd8b19c60a 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 @@ -30,8 +30,10 @@ import org.eclipse.leshan.core.node.LwM2mResource; import org.eclipse.leshan.core.observation.Observation; 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.response.ObserveResponse; +import org.eclipse.leshan.core.response.ReadCompositeResponse; import org.eclipse.leshan.core.response.ReadResponse; import org.eclipse.leshan.server.registration.Registration; import org.springframework.context.annotation.Lazy; @@ -89,6 +91,7 @@ import javax.annotation.PreDestroy; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -315,6 +318,24 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl } } + public void onUpdateValueAfterReadCompositeResponse(Registration registration, ReadCompositeResponse response) { + log.warn("201) ReadCompositeResponse: [{}]", response); + if (response.getContent() != null) { + LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint()); + response.getContent().forEach((k, v) -> { + if (v != null) { + if (v instanceof LwM2mObject) { + this.updateObjectResourceValue(lwM2MClient, (LwM2mObject) v, k.toString()); + } else if (v instanceof LwM2mObjectInstance) { + this.updateObjectInstanceResourceValue(lwM2MClient, (LwM2mObjectInstance) v, k.toString()); + } else if (v instanceof LwM2mResource) { + this.updateResourcesValue(lwM2MClient, (LwM2mResource) v, k.toString()); + } + } + }); + } + } + /** * @param sessionInfo - * @param deviceProfile - @@ -406,6 +427,17 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl if (supportedObjects != null && supportedObjects.size() > 0) { // #1 this.sendReadRequests(lwM2MClient, profile, supportedObjects); + // test composite + String[] paths = new String[]{"/3/0", "/1/0", "/5/0"}; +// String [] paths = new String[] {"/5"}; +// String [] paths = new String[] {"/"}; +// String [] paths = new String[] {"/9"}; +// defaultLwM2MDownlinkMsgHandler.sendReadCompositeRequest(lwM2MClient, paths, this); + Map nodes = new HashMap<>(); + nodes.put("/3/0/14", "+02"); + nodes.put("/1/0/2", 100); + nodes.put("/5/0/1", "coap://localhost:5685"); +// defaultLwM2MDownlinkMsgHandler.sendWriteCompositeRequest(lwM2MClient, nodes, this); this.sendObserveRequests(lwM2MClient, profile, supportedObjects); this.sendWriteAttributeRequests(lwM2MClient, profile, supportedObjects); // Removed. Used only for debug. @@ -682,6 +714,14 @@ public class DefaultLwM2MUplinkMsgHandler extends LwM2MExecutorAwareService impl } } + @Override + public void onWriteCompositeResponseOk(LwM2mClient client, WriteCompositeRequest request) { + log.warn("202) ReadCompositeResponse: [{}]", request.getNodes()); + request.getNodes().forEach((k, v) -> { + this.updateResourcesValue(client, (LwM2mResource) v, k.toString()); + }); + } + //TODO: review and optimize the logic to minimize number of the requests to device. private void onDeviceProfileUpdate(List clients, DeviceProfile deviceProfile) { var oldProfile = clientContext.getProfile(deviceProfile.getUuidId()); 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 ca01d197e8..b6fdf56a33 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 @@ -16,7 +16,9 @@ package org.thingsboard.server.transport.lwm2m.server.uplink; import org.eclipse.leshan.core.observation.Observation; +import org.eclipse.leshan.core.request.WriteCompositeRequest; import org.eclipse.leshan.core.request.WriteRequest; +import org.eclipse.leshan.core.response.ReadCompositeResponse; import org.eclipse.leshan.core.response.ReadResponse; import org.eclipse.leshan.server.registration.Registration; import org.thingsboard.server.common.data.Device; @@ -40,6 +42,8 @@ public interface LwM2mUplinkMsgHandler { void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response); + void onUpdateValueAfterReadCompositeResponse(Registration registration, ReadCompositeResponse response); + void onDeviceProfileUpdate(TransportProtos.SessionInfoProto sessionInfo, DeviceProfile deviceProfile); void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional deviceProfileOpt); @@ -52,6 +56,8 @@ public interface LwM2mUplinkMsgHandler { void onWriteResponseOk(LwM2mClient client, String path, WriteRequest request); + void onWriteCompositeResponseOk(LwM2mClient client, WriteCompositeRequest request); + void onToTransportUpdateCredentials(TransportProtos.SessionInfoProto sessionInfo, TransportProtos.ToTransportUpdateCredentialsProto updateCredentials); LwM2MTransportServerConfig getConfig();