Browse Source

Refactoring to extract AttributeService

pull/4752/head
Andrii Shvaika 5 years ago
parent
commit
dc7c96c4f0
  1. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/bootstrap/secure/LwM2MBootstrapSecurityStore.java
  2. 3
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mServerListener.java
  3. 6
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/LwM2mSessionMsgListener.java
  4. 150
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java
  5. 32
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/LwM2MAttributesService.java
  6. 6
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClient.java
  7. 8
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mFwSwUpdate.java
  8. 196
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java
  9. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveCallback.java
  10. 2
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadCallback.java
  11. 83
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java
  12. 19
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaState.java
  13. 24
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java
  14. 37
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/DefaultLwM2MRpcRequestHandler.java
  15. 284
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.java
  16. 210
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java
  17. 9
      common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java

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

@ -155,7 +155,7 @@ public class LwM2MBootstrapSecurityStore implements BootstrapSecurityStore {
LwM2MServerBootstrap profileLwm2mServer = JacksonUtil.fromString(JacksonUtil.toString(bootstrapObject.getLwm2mServer()), LwM2MServerBootstrap.class);
UUID sessionUUiD = UUID.randomUUID();
TransportProtos.SessionInfoProto sessionInfo = helper.getValidateSessionInfo(store.getMsg(), sessionUUiD.getMostSignificantBits(), sessionUUiD.getLeastSignificantBits());
context.getTransportService().registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(null, null, sessionInfo));
context.getTransportService().registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(null, null, null, sessionInfo));
if (this.getValidatedSecurityMode(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap, lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer)) {
lwM2MBootstrapConfig.bootstrapServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.bootstrapServer, profileServerBootstrap);
lwM2MBootstrapConfig.lwm2mServer = new LwM2MServerBootstrap(lwM2MBootstrapConfig.lwm2mServer, profileLwm2mServer);

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

@ -93,8 +93,7 @@ public class LwM2mServerListener {
@Override
public void onResponse(Observation observation, Registration registration, ObserveResponse response) {
if (registration != null) {
service.onUpdateValueAfterReadResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(),
registration), response, null);
service.onUpdateValueAfterReadResponse(registration, convertPathFromObjectIdToIdVer(observation.getPath().toString(), registration), response);
}
}

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

@ -31,6 +31,7 @@ import org.thingsboard.server.gen.transport.TransportProtos.SessionCloseNotifica
import org.thingsboard.server.gen.transport.TransportProtos.ToDeviceRpcRequestMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToServerRpcResponseMsg;
import org.thingsboard.server.gen.transport.TransportProtos.ToTransportUpdateCredentialsProto;
import org.thingsboard.server.transport.lwm2m.server.attributes.LwM2MAttributesService;
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcRequestHandler;
import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler;
import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler;
@ -42,17 +43,18 @@ import java.util.UUID;
@RequiredArgsConstructor
public class LwM2mSessionMsgListener implements GenericFutureListener<Future<? super Void>>, SessionMsgListener {
private final LwM2mUplinkMsgHandler handler;
private final LwM2MAttributesService attributesService;
private final LwM2MRpcRequestHandler rpcHandler;
private final TransportProtos.SessionInfoProto sessionInfo;
@Override
public void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse) {
this.handler.onGetAttributesResponse(getAttributesResponse, this.sessionInfo);
this.attributesService.onGetAttributesResponse(getAttributesResponse, this.sessionInfo);
}
@Override
public void onAttributeUpdate(AttributeUpdateNotificationMsg attributeUpdateNotification) {
this.handler.onAttributeUpdate(attributeUpdateNotification, this.sessionInfo);
this.attributesService.onAttributeUpdate(attributeUpdateNotification, this.sessionInfo);
}
@Override

150
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/attributes/DefaultLwM2MAttributesService.java

@ -0,0 +1,150 @@
/**
* 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.attributes;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.model.ResourceModel;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.ota.OtaPackageKey;
import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.common.data.ota.OtaPackageUtil;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.GetAttributeResponseMsg;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import java.util.Collection;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getValueFromKvProto;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_ERROR;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.isFwSwWords;
@Slf4j
@Service
@TbLwM2mTransportComponent
@RequiredArgsConstructor
public class DefaultLwM2MAttributesService implements LwM2MAttributesService {
//TODO: add timeout logic
private final AtomicInteger reqIdSeq = new AtomicInteger();
private final Map<Integer, SettableFuture<List<TransportProtos.TsKvProto>>> futures;
private final TransportService transportService;
@Override
public ListenableFuture<List<TransportProtos.TsKvProto>> getSharedAttributes(LwM2mClient client, Collection<String> keys) {
SettableFuture<List<TransportProtos.TsKvProto>> future = SettableFuture.create();
int requestId = reqIdSeq.incrementAndGet();
futures.put(requestId, future);
transportService.process(client.getSession(), TransportProtos.GetAttributeRequestMsg.newBuilder().setRequestId(requestId).
addAllSharedAttributeNames(keys).build(), new TransportServiceCallback<Void>() {
@Override
public void onSuccess(Void msg) {
}
@Override
public void onError(Throwable e) {
SettableFuture<List<TransportProtos.TsKvProto>> callback = futures.remove(requestId);
if (callback != null) {
callback.setException(e);
}
}
});
return future;
}
@Override
public void onGetAttributesResponse(GetAttributeResponseMsg getAttributesResponse, TransportProtos.SessionInfoProto sessionInfo) {
var callback = futures.remove(getAttributesResponse.getRequestId());
if (callback != null) {
callback.set(getAttributesResponse.getSharedAttributeListList());
}
}
/**
* Update - send request in change value resources in Client
* 1. FirmwareUpdate:
* - If msg.getSharedUpdatedList().forEach(tsKvProto -> {tsKvProto.getKv().getKey().indexOf(FIRMWARE_UPDATE_PREFIX, 0) == 0
* 2. Shared Other AttributeUpdate
* -- Path to resources from profile equal keyName or from ModelObject equal name
* -- Only for resources: isWritable && isPresent as attribute in profile -> LwM2MClientProfile (format: CamelCase)
* 3. Delete - nothing
*
* @param msg -
*/
@Override
public void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo) {
// LwM2mClient lwM2MClient = clientContext.getClientBySessionInfo(sessionInfo);
// if (msg.getSharedUpdatedCount() > 0 && lwM2MClient != null) {
// log.warn("2) OnAttributeUpdate, SharedUpdatedList() [{}]", msg.getSharedUpdatedList());
// msg.getSharedUpdatedList().forEach(tsKvProto -> {
// String pathName = tsKvProto.getKv().getKey();
// String pathIdVer = this.getObjectIdByKeyNameFromProfile(sessionInfo, pathName);
// Object valueNew = getValueFromKvProto(tsKvProto.getKv());
// if ((OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName)
// && (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion())))
// || (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.TITLE).equals(pathName)
// && (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentTitle())))) {
// this.getInfoFirmwareUpdate(lwM2MClient, null);
// } else if ((OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.VERSION).equals(pathName)
// && (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentVersion())))
// || (OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.TITLE).equals(pathName)
// && (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentTitle())))) {
// this.getInfoSoftwareUpdate(lwM2MClient, null);
// }
// if (pathIdVer != null) {
// ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.config
// .getModelProvider());
// if (resourceModel != null && resourceModel.operations.isWritable()) {
// this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), valueNew, pathIdVer);
// } else {
// log.error("Resource path - [{}] value - [{}] is not Writable and cannot be updated", pathIdVer, valueNew);
// String logMsg = String.format("%s: attributeUpdate: Resource path - %s value - %s is not Writable and cannot be updated",
// LOG_LWM2M_ERROR, pathIdVer, valueNew);
// this.logToTelemetry(lwM2MClient, logMsg);
// }
// } else if (!isFwSwWords(pathName)) {
// log.error("Resource name name - [{}] value - [{}] is not present as attribute/telemetry in profile and cannot be updated", pathName, valueNew);
// String logMsg = String.format("%s: attributeUpdate: attribute name - %s value - %s is not present as attribute in profile and cannot be updated",
// LOG_LWM2M_ERROR, pathName, valueNew);
// this.logToTelemetry(lwM2MClient, logMsg);
// }
//
// });
// } else if (msg.getSharedDeletedCount() > 0 && lwM2MClient != null) {
// msg.getSharedUpdatedList().forEach(tsKvProto -> {
// String pathName = tsKvProto.getKv().getKey();
// Object valueNew = getValueFromKvProto(tsKvProto.getKv());
// if (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName) && !valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion())) {
// lwM2MClient.getFwUpdate().setCurrentVersion((String) valueNew);
// }
// });
// log.info("[{}] delete [{}] onAttributeUpdate", msg.getSharedDeletedList(), sessionInfo);
// } else if (lwM2MClient == null) {
// log.error("OnAttributeUpdate, lwM2MClient is null");
// }
}
}

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

@ -0,0 +1,32 @@
/**
* Copyright © 2016-2021 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.transport.lwm2m.server.attributes;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import java.util.Collection;
import java.util.List;
public interface LwM2MAttributesService {
ListenableFuture<List<TransportProtos.TsKvProto>> getSharedAttributes(LwM2mClient client, Collection<String> keys);
void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg getAttributesResponse, TransportProtos.SessionInfoProto sessionInfo);
void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg attributeUpdateNotification, TransportProtos.SessionInfoProto sessionInfo);
}

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

@ -75,7 +75,7 @@ public class LwM2mClient implements Cloneable {
@Getter
private final Map<String, ResourceValue> resources;
@Getter
private final Map<String, TsKvProto> delayedRequests;
private final Map<String, TsKvProto> sharedAttributes;
@Getter
private final List<String> pendingReadRequests;
@Getter
@ -121,7 +121,7 @@ public class LwM2mClient implements Cloneable {
this.nodeId = nodeId;
this.endpoint = endpoint;
this.lock = new ReentrantLock();
this.delayedRequests = new ConcurrentHashMap<>();
this.sharedAttributes = new ConcurrentHashMap<>();
this.pendingReadRequests = new CopyOnWriteArrayList<>();
this.resources = new ConcurrentHashMap<>();
this.queuedRequests = new ConcurrentLinkedQueue<>();
@ -371,7 +371,7 @@ public class LwM2mClient implements Cloneable {
}
if (this.pendingReadRequests.size() == 0) {
this.init = true;
serviceImpl.putDelayedUpdateResourcesThingsboard(this);
serviceImpl.initAttributes(this);
}
}

8
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mFwSwUpdate.java

@ -24,7 +24,6 @@ import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType;
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest;
import org.thingsboard.server.transport.lwm2m.server.uplink.DefaultLwM2MUplinkMsgHandler;
import org.thingsboard.server.transport.lwm2m.server.downlink.LwM2mDownlinkMsgHandler;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil;
@ -127,9 +126,6 @@ public class LwM2mFwSwUpdate {
private final List<String> pendingInfoRequestsStart;
@Getter
@Setter
private volatile LwM2mClientRpcRequest rpcRequest;
@Getter
@Setter
private volatile int updateStrategy;
public LwM2mFwSwUpdate(LwM2mUplinkMsgHandler handler, LwM2mClient lwM2MClient, OtaPackageType type, int updateStrategy) {
@ -214,10 +210,10 @@ public class LwM2mFwSwUpdate {
} else {
String msgError = "FirmWareId is null.";
log.warn("6) [{}]", msgError);
if (this.rpcRequest != null) {
// if (this.rpcRequest != null) {
// TODO: refactor.
// handler.sentRpcResponse(this.rpcRequest, CONTENT.name(), msgError, LOG_LW2M_ERROR);
}
// }
log.error(msgError);
this.sendLogs(handler, WRITE_REPLACE.name(), LOG_LWM2M_ERROR, msgError);
}

196
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/DefaultLwM2mDownlinkMsgHandler.java

@ -26,6 +26,7 @@ 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.CancelObservationRequest;
import org.eclipse.leshan.core.request.ContentFormat;
import org.eclipse.leshan.core.request.DeleteRequest;
import org.eclipse.leshan.core.request.DiscoverRequest;
@ -49,13 +50,10 @@ import org.eclipse.leshan.server.registration.Registration;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.server.common.data.device.data.lwm2m.ObjectAttributes;
import org.thingsboard.server.common.transport.util.JsonUtils;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
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.client.LwM2mClientContext;
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
import javax.annotation.PostConstruct;
@ -289,200 +287,8 @@ public class DefaultLwM2mDownlinkMsgHandler implements LwM2mDownlinkMsgHandler {
default:
throw new IllegalArgumentException("Not supported type:" + type.name());
}
// TODO: throw exception and execute callback.
//// if (rpcRequest != null) {
//// String patn = "/" + objectId + "/" + instanceId + "/" + resourceId;
//// String errorMsg = String.format("Bad ResourceModel Operations (E): Resource path - %s ResourceModel type - %s", patn, type);
//// rpcRequest.setErrorMsg(errorMsg);
//// }
// return null;
// } catch (NumberFormatException e) {
// String patn = "/" + objectId + "/" + instanceId + "/" + resourceId;
// String msg = String.format(LOG_LW2M_ERROR + ": NumberFormatException: Resource path - %s type - %s value - %s msg error - %s SendRequest to Client",
// patn, type, value, e.toString());
// handler.sendLogsToThingsboard(client, msg);
// log.error("Path: [{}] type: [{}] value: [{}] errorMsg: [{}]]", patn, type, value, e.toString());
// if (rpcRequest != null) {
// String errorMsg = String.format("NumberFormatException: Resource path - %s type - %s value - %s", patn, type, value);
// handler.sentRpcResponse(rpcRequest, BAD_REQUEST.getName(), errorMsg, LOG_LW2M_ERROR);
// }
// return null;
// }
}
private void handleResponse(LwM2mClient lwM2mClient, final String path, LwM2mResponse response,
SimpleDownlinkRequest request, LwM2mClientRpcRequest rpcRequest) {
responseRequestExecutor.submit(() -> {
try {
this.sendResponse(lwM2mClient, path, response, request, rpcRequest);
} catch (Exception e) {
log.error("[{}] endpoint [{}] path [{}] Exception Unable to after send response.", lwM2mClient.getRegistration().getEndpoint(), path, e);
}
});
}
/**
* processing a response from a client
*
* @param path -
* @param response -
*/
private void sendResponse(LwM2mClient lwM2mClient, String path, LwM2mResponse response,
SimpleDownlinkRequest request, LwM2mClientRpcRequest rpcRequest) {
// Registration registration = lwM2mClient.getRegistration();
// String pathIdVer = convertPathFromObjectIdToIdVer(path, registration);
// String msgLog = "";
// if (response instanceof ReadResponse) {
// handler.onUpdateValueAfterReadResponse(registration, pathIdVer, (ReadResponse) response, rpcRequest);
// } else if (response instanceof DeleteResponse) {
// log.warn("11) [{}] Path [{}] DeleteResponse", pathIdVer, response);
// if (rpcRequest != null) {
// rpcRequest.setInfoMsg(null);
// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), null, null);
// }
// } else if (response instanceof DiscoverResponse) {
// String discoverValue = Link.serialize(((DiscoverResponse) response).getObjectLinks());
// msgLog = String.format("%s: type operation: %s path: %s value: %s",
// LOG_LW2M_INFO, DISCOVER.name(), request.getPath().toString(), discoverValue);
// handler.sendLogsToThingsboard(lwM2mClient, msgLog);
// log.warn("DiscoverResponse: [{}]", (DiscoverResponse) response);
// if (rpcRequest != null) {
// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), discoverValue, LOG_LW2M_VALUE);
// }
// } else if (response instanceof ExecuteResponse) {
// msgLog = String.format("%s: type operation: %s path: %s",
// LOG_LW2M_INFO, EXECUTE.name(), request.getPath().toString());
// log.warn("9) [{}] ", msgLog);
// handler.sendLogsToThingsboard(lwM2mClient, msgLog);
// if (rpcRequest != null) {
// msgLog = String.format("Start %s path: %S. Preparation finished: %s", EXECUTE.name(), path, rpcRequest.getInfoMsg());
// rpcRequest.setInfoMsg(msgLog);
// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), path, LOG_LW2M_INFO);
// }
//
// } else if (response instanceof WriteAttributesResponse) {
// msgLog = String.format("%s: type operation: %s path: %s value: %s",
// LOG_LW2M_INFO, WRITE_ATTRIBUTES.name(), request.getPath().toString(), ((WriteAttributesRequest) request).getAttributes().toString());
// handler.sendLogsToThingsboard(lwM2mClient, msgLog);
// log.warn("12) [{}] Path [{}] WriteAttributesResponse", pathIdVer, response);
// if (rpcRequest != null) {
// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), response.toString(), LOG_LW2M_VALUE);
// }
// } else if (response instanceof WriteResponse) {
// msgLog = String.format("Type operation: Write path: %s", pathIdVer);
// log.warn("10) [{}] response: [{}]", msgLog, response);
// this.infoWriteResponse(lwM2mClient, response, request, rpcRequest);
// handler.onWriteResponseOk(registration, pathIdVer, (WriteRequest) request);
// }
}
// private void infoWriteResponse(LwM2mClient lwM2mClient, LwM2mResponse response, SimpleDownlinkRequest
// request, LwM2mClientRpcRequest rpcRequest) {
// try {
// Registration registration = lwM2mClient.getRegistration();
// LwM2mNode node = ((WriteRequest) request).getNode();
// String msg = null;
// Object value;
// if (node instanceof LwM2mObject) {
// msg = String.format("%s: Update finished successfully: Lwm2m code - %d Source path: %s value: %s",
// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), ((LwM2mObject) node).toString());
// } else if (node instanceof LwM2mObjectInstance) {
// msg = String.format("%s: Update finished successfully: Lwm2m code - %d Source path: %s value: %s",
// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), ((LwM2mObjectInstance) node).prettyPrint());
// } else if (node instanceof LwM2mSingleResource) {
// LwM2mSingleResource singleResource = (LwM2mSingleResource) node;
// if (singleResource.getType() == ResourceModel.Type.STRING || singleResource.getType() == ResourceModel.Type.OPAQUE) {
// int valueLength;
// if (singleResource.getType() == ResourceModel.Type.STRING) {
// valueLength = ((String) singleResource.getValue()).length();
// value = ((String) singleResource.getValue())
// .substring(Math.min(valueLength, config.getLogMaxLength())).trim();
//
// } else {
// valueLength = ((byte[]) singleResource.getValue()).length;
// value = new String(Arrays.copyOf(((byte[]) singleResource.getValue()),
// Math.min(valueLength, config.getLogMaxLength()))).trim();
// }
// value = valueLength > config.getLogMaxLength() ? value + "..." : value;
// msg = String.format("%s: Update finished successfully: Lwm2m code - %d Resource path: %s length: %s value: %s",
// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), valueLength, value);
// } else {
// value = this.converter.convertValue(singleResource.getValue(),
// singleResource.getType(), ResourceModel.Type.STRING, request.getPath());
// msg = String.format("%s: Update finished successfully. Lwm2m code: %d Resource path: %s value: %s",
// LOG_LW2M_INFO, response.getCode().getCode(), request.getPath().toString(), value);
// }
// }
// if (msg != null) {
// handler.sendLogsToThingsboard(lwM2mClient, msg);
// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) || request.getPath().toString().equals(SW_PACKAGE_ID)) {
// this.afterWriteSuccessFwSwUpdate(registration, request);
// if (rpcRequest != null) {
// rpcRequest.setInfoMsg(msg);
// }
// } else if (rpcRequest != null) {
// handler.sentRpcResponse(rpcRequest, response.getCode().getName(), msg, LOG_LW2M_INFO);
// }
// }
// } catch (Exception e) {
// log.trace("Fail convert value from request to string. ", e);
// }
// }
/**
* After finish operation FwSwUpdate Write (success):
* fw_state/sw_state = DOWNLOADED
* send operation Execute
*/
// private void afterWriteSuccessFwSwUpdate(Registration registration, SimpleDownlinkRequest request) {
// LwM2mClient client = this.lwM2mClientContext.getClientByRegistrationId(registration.getId());
// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) && client.getFwUpdate() != null) {
// client.getFwUpdate().setStateUpdate(DOWNLOADED.name());
// client.getFwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_INFO, null);
// }
// if (request.getPath().toString().equals(SW_PACKAGE_ID) && client.getSwUpdate() != null) {
// client.getSwUpdate().setStateUpdate(DOWNLOADED.name());
// client.getSwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_INFO, null);
// }
// }
/**
* After finish operation FwSwUpdate Write (error): fw_state = FAILED
*/
// private void afterWriteFwSWUpdateError(Registration registration, SimpleDownlinkRequest request, String
// msgError) {
// LwM2mClient client = this.lwM2mClientContext.getClientByRegistrationId(registration.getId());
// if (request.getPath().toString().equals(FW_PACKAGE_5_ID) && client.getFwUpdate() != null) {
// client.getFwUpdate().setStateUpdate(FAILED.name());
// client.getFwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_ERROR, msgError);
// }
// if (request.getPath().toString().equals(SW_PACKAGE_ID) && client.getSwUpdate() != null) {
// client.getSwUpdate().setStateUpdate(FAILED.name());
// client.getSwUpdate().sendLogs(this.handler, WRITE_REPLACE.name(), LOG_LW2M_ERROR, msgError);
// }
// }
// private void afterExecuteFwSwUpdateError(Registration registration, SimpleDownlinkRequest request, String
// msgError) {
// LwM2mClient client = this.lwM2mClientContext.getClientByRegistrationId(registration.getId());
// if (request.getPath().toString().equals(FW_UPDATE_ID) && client.getFwUpdate() != null) {
// client.getFwUpdate().sendLogs(this.handler, EXECUTE.name(), LOG_LW2M_ERROR, msgError);
// }
// if (request.getPath().toString().equals(SW_INSTALL_ID) && client.getSwUpdate() != null) {
// client.getSwUpdate().sendLogs(this.handler, EXECUTE.name(), LOG_LW2M_ERROR, msgError);
// }
// }
// private void afterObserveCancel(LwM2mClient lwM2mClient, int observeCancelCnt, String
// observeCancelMsg, LwM2mClientRpcRequest rpcRequest) {
// handler.sendLogsToThingsboard(lwM2mClient, observeCancelMsg);
// log.warn("[{}]", observeCancelMsg);
// if (rpcRequest != null) {
// rpcRequest.setInfoMsg(String.format("Count: %d", observeCancelCnt));
// handler.sentRpcResponse(rpcRequest, CONTENT.name(), null, LOG_LW2M_INFO);
// }
// }
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!");

2
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MObserveCallback.java

@ -31,6 +31,6 @@ public class TbLwM2MObserveCallback extends TbLwM2MTargetedCallback<ObserveReque
@Override
public void onSuccess(ObserveRequest request, ObserveResponse response) {
super.onSuccess(request, response);
handler.onUpdateValueAfterReadResponse(client.getRegistration(), versionedId, response, null);
handler.onUpdateValueAfterReadResponse(client.getRegistration(), versionedId, response);
}
}

2
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/downlink/TbLwM2MReadCallback.java

@ -31,7 +31,7 @@ public class TbLwM2MReadCallback extends TbLwM2MTargetedCallback<ReadRequest, Re
@Override
public void onSuccess(ReadRequest request, ReadResponse response) {
super.onSuccess(request, response);
handler.onUpdateValueAfterReadResponse(client.getRegistration(), versionedId, response, null);
handler.onUpdateValueAfterReadResponse(client.getRegistration(), versionedId, response);
}
}

83
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/DefaultLwM2MOtaUpdateService.java

@ -0,0 +1,83 @@
/**
* 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.ota;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.ota.OtaPackageKey;
import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
import org.thingsboard.server.transport.lwm2m.server.attributes.LwM2MAttributesService;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClientContext;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import static org.thingsboard.server.common.data.ota.OtaPackageUtil.getAttributeKey;
@Slf4j
@Service
@TbLwM2mTransportComponent
@RequiredArgsConstructor
public class DefaultLwM2MOtaUpdateService implements LwM2MOtaUpdateService {
private static final String FW_NAME_ID = "/5/0/6";
private static final String FW_VER_ID = "/5/0/7";
private static final String SW_NAME_ID = "/9/0/0";
private static final String SW_VER_ID = "/9/0/1";
private final Map<String, LwM2MClientOtaState> fwStates = new ConcurrentHashMap<>();
private final Map<String, LwM2MClientOtaState> swStates = new ConcurrentHashMap<>();
private final LwM2MAttributesService attributesService;
private final TransportService transportService;
private final LwM2mClientContext clientContext;
@Override
public void init(LwM2mClient client) {
//TODO: check that the client supports FW and SW by checking the supported objects in the model.
List<String> attributesToFetch = new ArrayList<>();
if (client.isValidObjectVersion(FW_NAME_ID) || client.isValidObjectVersion(FW_VER_ID)) {
LwM2MClientOtaState fwState = getOrInitFwSate(client);
attributesToFetch.add(getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.TITLE));
attributesToFetch.add(getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION));
}
if (client.isValidObjectVersion(SW_NAME_ID) || client.isValidObjectVersion(SW_VER_ID)) {
LwM2MClientOtaState swState = getOrInitSwSate(client);
attributesToFetch.add(getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.TITLE));
attributesToFetch.add(getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.VERSION));
}
var future = attributesService.getSharedAttributes(client, attributesToFetch);
}
private LwM2MClientOtaState getOrInitFwSate(LwM2mClient client) {
//TODO: fetch state from the cache.
return fwStates.computeIfAbsent(client.getEndpoint(), endpoint -> new LwM2MClientOtaState());
}
private LwM2MClientOtaState getOrInitSwSate(LwM2mClient client) {
//TODO: fetch state from the cache.
return swStates.computeIfAbsent(client.getEndpoint(), endpoint -> new LwM2MClientOtaState());
}
}

19
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MClientOtaState.java

@ -0,0 +1,19 @@
/**
* 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.ota;
public class LwM2MClientOtaState {
}

24
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/ota/LwM2MOtaUpdateService.java

@ -0,0 +1,24 @@
/**
* 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.ota;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
public interface LwM2MOtaUpdateService {
void init(LwM2mClient client);
}

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

@ -58,11 +58,6 @@ import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.stream.Collectors;
import static org.eclipse.californium.core.coap.CoAP.ResponseCode.BAD_REQUEST;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_ERROR;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_INFO;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.LOG_LWM2M_VALUE;
@Slf4j
@Service
@TbLwM2mTransportComponent
@ -258,17 +253,6 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler {
transportService.process(sessionInfo, msg, null);
}
private void sendErrorRpcResponse(LwM2mClientRpcRequest lwm2mClientRpcRequest, String msgError, TransportProtos.SessionInfoProto sessionInfo) {
if (lwm2mClientRpcRequest == null) {
lwm2mClientRpcRequest = new LwM2mClientRpcRequest();
}
lwm2mClientRpcRequest.setResponseCode(BAD_REQUEST.name());
if (lwm2mClientRpcRequest.getErrorMsg() == null) {
lwm2mClientRpcRequest.setErrorMsg(msgError);
}
this.onToDeviceRpcResponse(lwm2mClientRpcRequest.getDeviceRpcResponseResultMsg(), sessionInfo);
}
private void cleanupOldSessions() {
log.warn("4.1) before rpcSubscriptions.size(): [{}]", rpcSubscriptions.size());
if (rpcSubscriptions.size() > 0) {
@ -281,27 +265,6 @@ public class DefaultLwM2MRpcRequestHandler implements LwM2MRpcRequestHandler {
log.warn("4.4) after rpcSubscriptions.size(): [{}]", rpcSubscriptions.size());
}
public void sentRpcResponse(LwM2mClientRpcRequest rpcRequest, String requestCode, String msg, String typeMsg) {
rpcRequest.setResponseCode(requestCode);
if (LOG_LWM2M_ERROR.equals(typeMsg)) {
rpcRequest.setInfoMsg(null);
rpcRequest.setValueMsg(null);
if (rpcRequest.getErrorMsg() == null) {
msg = msg.isEmpty() ? null : msg;
rpcRequest.setErrorMsg(msg);
}
} else if (LOG_LWM2M_INFO.equals(typeMsg)) {
if (rpcRequest.getInfoMsg() == null) {
rpcRequest.setInfoMsg(msg);
}
} else if (LOG_LWM2M_VALUE.equals(typeMsg)) {
if (rpcRequest.getValueMsg() == null) {
rpcRequest.setValueMsg(msg);
}
}
this.onToDeviceRpcResponse(rpcRequest.getDeviceRpcResponseResultMsg(), rpcRequest.getSessionInfo());
}
@Override
public void onToDeviceRpcResponse(TransportProtos.ToDeviceRpcResponseMsg toDeviceResponse, TransportProtos.SessionInfoProto sessionInfo) {
log.warn("5) onToDeviceRpcResponse: [{}], sessionUUID: [{}]", toDeviceResponse, new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB()));

284
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/rpc/LwM2mClientRpcRequest.java

@ -1,284 +0,0 @@
/**
* 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;
import com.google.gson.Gson;
import com.google.gson.JsonObject;
import com.google.gson.reflect.TypeToken;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.eclipse.leshan.core.node.LwM2mPath;
import org.eclipse.leshan.server.registration.Registration;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.server.uplink.LwM2mUplinkMsgHandler;
import java.util.Map;
import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeoutException;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.ERROR_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FINISH_JSON_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FINISH_VALUE_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.INFO_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.KEY_NAME_KEY;
import org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.DISCOVER_ALL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.EXECUTE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.FW_UPDATE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.OBSERVE_CANCEL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.OBSERVE_READ_ALL;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.WRITE_ATTRIBUTES;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.WRITE_REPLACE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mOperationType.WRITE_UPDATE;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.METHOD_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.PARAMS_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.RESULT_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.SEPARATOR_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.START_JSON_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.TARGET_ID_VER_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.VALUE_KEY;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validPathIdVer;
@Slf4j
@Data
public class LwM2mClientRpcRequest {
private Registration registration;
private TransportProtos.SessionInfoProto sessionInfo;
private String bodyParams;
private int requestId;
private LwM2mOperationType typeOper;
private String key;
private String targetIdVer;
private Object value;
private Map<String, Object> params;
private String errorMsg;
private String valueMsg;
private String infoMsg;
private String responseCode;
public LwM2mClientRpcRequest() {
}
public LwM2mClientRpcRequest(LwM2mOperationType lwM2mTypeOper, String bodyParams, int requestId,
TransportProtos.SessionInfoProto sessionInfo, Registration registration, LwM2mUplinkMsgHandler handler) {
this.registration = registration;
this.sessionInfo = sessionInfo;
this.requestId = requestId;
if (lwM2mTypeOper != null) {
this.typeOper = lwM2mTypeOper;
} else {
this.errorMsg = METHOD_KEY + " - " + typeOper + " is not valid.";
}
if (this.errorMsg == null && !bodyParams.equals("null")) {
this.bodyParams = bodyParams;
this.init(handler);
}
}
public TransportProtos.ToDeviceRpcResponseMsg getDeviceRpcResponseResultMsg() {
JsonObject payloadResp = new JsonObject();
payloadResp.addProperty(RESULT_KEY, this.responseCode);
if (this.errorMsg != null) {
payloadResp.addProperty(ERROR_KEY, this.errorMsg);
} else if (this.valueMsg != null) {
payloadResp.addProperty(VALUE_KEY, this.valueMsg);
} else if (this.infoMsg != null) {
payloadResp.addProperty(INFO_KEY, this.infoMsg);
}
return TransportProtos.ToDeviceRpcResponseMsg.newBuilder()
.setPayload(payloadResp.getAsJsonObject().toString())
.setRequestId(this.requestId)
.build();
}
private void init(LwM2mUplinkMsgHandler handler) {
try {
// #1
if (this.bodyParams.contains(KEY_NAME_KEY)) {
String targetIdVerStr = this.getValueKeyFromBody(KEY_NAME_KEY);
if (targetIdVerStr != null) {
String targetIdVer = handler.getObjectIdByKeyNameFromProfile(sessionInfo, targetIdVerStr);
if (targetIdVer != null) {
this.targetIdVer = targetIdVer;
this.setInfoMsg(String.format("Changed by: key - %s, pathIdVer - %s",
targetIdVerStr, targetIdVer));
}
}
}
if (this.getTargetIdVer() == null && this.bodyParams.contains(TARGET_ID_VER_KEY)) {
this.setValidTargetIdVerKey();
}
if (this.bodyParams.contains(VALUE_KEY)) {
this.value = this.getValueKeyFromBody(VALUE_KEY);
}
try {
if (this.bodyParams.contains(PARAMS_KEY)) {
this.setValidParamsKey(handler);
}
} catch (Exception e) {
this.setErrorMsg(String.format("Params of request is bad Json format. %s", e.getMessage()));
}
if (this.getTargetIdVer() == null
&& !(OBSERVE_READ_ALL == this.getTypeOper()
|| DISCOVER_ALL == this.getTypeOper()
|| OBSERVE_CANCEL == this.getTypeOper()
|| FW_UPDATE == this.getTypeOper())) {
this.setErrorMsg(TARGET_ID_VER_KEY + " and " +
KEY_NAME_KEY + " is null or bad format");
}
/**
* EXECUTE && WRITE_REPLACE - only for Resource or ResourceInstance
*/
else if (this.getTargetIdVer() != null
&& (EXECUTE == this.getTypeOper()
|| WRITE_REPLACE == this.getTypeOper())
&& !(new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(this.getTargetIdVer()))).isResource()
|| new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(this.getTargetIdVer()))).isResourceInstance())) {
this.setErrorMsg("Invalid parameter " + TARGET_ID_VER_KEY
+ ". Only Resource or ResourceInstance can be this operation");
}
} catch (Exception e) {
this.setErrorMsg(String.format("Bad format request. %s", e.getMessage()));
}
}
private void setValidTargetIdVerKey() {
String targetIdVerStr = this.getValueKeyFromBody(TARGET_ID_VER_KEY);
// targetIdVer without ver - ok
try {
// targetIdVer with/without ver - ok
this.targetIdVer = validPathIdVer(targetIdVerStr, this.registration);
if (this.targetIdVer != null) {
this.infoMsg = String.format("Changed by: pathIdVer - %s", this.targetIdVer);
}
} catch (Exception e) {
if (this.targetIdVer == null) {
this.errorMsg = TARGET_ID_VER_KEY + " - " + targetIdVerStr + " is not valid.";
}
}
}
private void setValidParamsKey(LwM2mUplinkMsgHandler handler) {
String paramsStr = this.getValueKeyFromBody(PARAMS_KEY);
if (paramsStr != null) {
String params2Json =
START_JSON_KEY
+ "\""
+ paramsStr
.replaceAll(SEPARATOR_KEY, "\"" + SEPARATOR_KEY + "\"")
.replaceAll(FINISH_VALUE_KEY, "\"" + FINISH_VALUE_KEY + "\"")
+ "\""
+ FINISH_JSON_KEY;
// jsonObject
Map<String, Object> params = new Gson().fromJson(params2Json, new TypeToken<ConcurrentHashMap<String, Object>>() {
}.getType());
if (WRITE_UPDATE == this.getTypeOper()) {
if (this.targetIdVer != null) {
Map<String, Object> paramsResourceId = this.convertParamsToResourceId((ConcurrentHashMap<String, Object>) params, handler);
if (paramsResourceId.size() > 0) {
this.setParams(paramsResourceId);
}
}
} else if (WRITE_ATTRIBUTES == this.getTypeOper()) {
this.setParams(params);
}
}
}
private String getValueKeyFromBody(String key) {
String valueKey = null;
int startInd = -1;
int finishInd = -1;
try {
switch (key) {
case KEY_NAME_KEY:
case TARGET_ID_VER_KEY:
case VALUE_KEY:
startInd = this.bodyParams.indexOf(SEPARATOR_KEY, this.bodyParams.indexOf(key));
finishInd = this.bodyParams.indexOf(FINISH_VALUE_KEY, this.bodyParams.indexOf(key));
if (startInd >= 0 && finishInd < 0) {
finishInd = this.bodyParams.indexOf(FINISH_JSON_KEY, this.bodyParams.indexOf(key));
}
break;
case PARAMS_KEY:
startInd = this.bodyParams.indexOf(START_JSON_KEY, this.bodyParams.indexOf(key));
finishInd = this.bodyParams.indexOf(FINISH_JSON_KEY, this.bodyParams.indexOf(key));
}
if (startInd >= 0 && finishInd > 0) {
valueKey = this.bodyParams.substring(startInd + 1, finishInd);
}
} catch (Exception e) {
log.error("", new TimeoutException());
}
/**
* ReplaceAll "\""
*/
if (StringUtils.trimToNull(valueKey) != null) {
char[] chars = valueKey.toCharArray();
for (int i = 0; i < chars.length; i++) {
if (chars[i] == 92 || chars[i] == 34) chars[i] = 32;
}
return key.equals(PARAMS_KEY) ? String.valueOf(chars) : String.valueOf(chars).replaceAll(" ", "");
}
return null;
}
private ConcurrentHashMap<String, Object> convertParamsToResourceId(ConcurrentHashMap<String, Object> params,
LwM2mUplinkMsgHandler serviceImpl) {
Map<String, Object> paramsIdVer = new ConcurrentHashMap<>();
LwM2mPath targetId = new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(this.targetIdVer)));
if (targetId.isObjectInstance()) {
params.forEach((k, v) -> {
try {
int id = Integer.parseInt(k);
paramsIdVer.put(String.valueOf(id), v);
} catch (NumberFormatException e) {
String targetIdVer = serviceImpl.getObjectIdByKeyNameFromProfile(sessionInfo, k);
if (targetIdVer != null) {
LwM2mPath lwM2mPath = new LwM2mPath(Objects.requireNonNull(fromVersionedIdToObjectId(targetIdVer)));
paramsIdVer.put(String.valueOf(lwM2mPath.getResourceId()), v);
}
/** WRITE_UPDATE*/
else {
String rezId = this.getRezIdByResourceNameAndObjectInstanceId(k, serviceImpl);
if (rezId != null) {
paramsIdVer.put(rezId, v);
}
}
}
});
}
return (ConcurrentHashMap<String, Object>) paramsIdVer;
}
private String getRezIdByResourceNameAndObjectInstanceId(String resourceName, LwM2mUplinkMsgHandler handler) {
// LwM2mClient lwM2mClient = handler.clientContext.getClientBySessionInfo(this.sessionInfo);
// return lwM2mClient != null ?
// lwM2mClient.getRezIdByResourceNameAndObjectInstanceId(resourceName, this.targetIdVer, handler.config.getModelProvider()) :
// null;
return null;
}
}

210
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/DefaultLwM2MUplinkMsgHandler.java

@ -33,6 +33,7 @@ import org.eclipse.leshan.core.response.ReadResponse;
import org.eclipse.leshan.server.registration.Registration;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.thingsboard.common.util.DonAsynchron;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.cache.ota.OtaPackageDataCache;
import org.thingsboard.server.common.data.Device;
@ -41,15 +42,12 @@ import org.thingsboard.server.common.data.device.data.lwm2m.ObjectAttributes;
import org.thingsboard.server.common.data.device.data.lwm2m.TelemetryMappingConfiguration;
import org.thingsboard.server.common.data.device.profile.Lwm2mDeviceProfileTransportConfiguration;
import org.thingsboard.server.common.data.id.OtaPackageId;
import org.thingsboard.server.common.data.ota.OtaPackageKey;
import org.thingsboard.server.common.data.ota.OtaPackageType;
import org.thingsboard.server.common.data.ota.OtaPackageUtil;
import org.thingsboard.server.common.transport.TransportService;
import org.thingsboard.server.common.transport.TransportServiceCallback;
import org.thingsboard.server.common.transport.adaptor.AdaptorException;
import org.thingsboard.server.common.transport.service.DefaultTransportService;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.AttributeUpdateNotificationMsg;
import org.thingsboard.server.gen.transport.TransportProtos.SessionEvent;
import org.thingsboard.server.gen.transport.TransportProtos.SessionInfoProto;
import org.thingsboard.server.queue.util.TbLwM2mTransportComponent;
@ -61,6 +59,7 @@ import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportContext;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil;
import org.thingsboard.server.transport.lwm2m.server.adaptors.LwM2MJsonAdaptor;
import org.thingsboard.server.transport.lwm2m.server.attributes.LwM2MAttributesService;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientState;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2MClientStateException;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
@ -83,7 +82,6 @@ import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteAttrib
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteResponseCallback;
import org.thingsboard.server.transport.lwm2m.server.downlink.TbLwM2MWriteReplaceRequest;
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2MRpcRequestHandler;
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest;
import org.thingsboard.server.transport.lwm2m.server.store.TbLwM2MDtlsSessionStore;
import org.thingsboard.server.transport.lwm2m.utils.LwM2mValueConverterImpl;
@ -107,7 +105,6 @@ import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPA
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.FAILED;
import static org.thingsboard.server.common.data.ota.OtaPackageUpdateStatus.INITIATED;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServerHelper.getValueFromKvProto;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.DEVICE_ATTRIBUTES_REQUEST;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_5_ID;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_RESULT_ID;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.FW_STATE_ID;
@ -119,8 +116,6 @@ import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.S
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertOtaUpdateValueToString;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.convertPathFromObjectIdToIdVer;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.fromVersionedIdToObjectId;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.getAckCallback;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.isFwSwWords;
import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportUtil.validateObjectVerFromKey;
@ -136,6 +131,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
private final TransportService transportService;
private final LwM2mTransportContext context;
private final LwM2MAttributesService attributesService;
public final LwM2MTransportServerConfig config;
public final OtaPackageDataCache otaPackageDataCache;
public final LwM2mTransportServerHelper helper;
@ -147,13 +143,14 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
public final Map<String, Integer> firmwareUpdateState;
public DefaultLwM2MUplinkMsgHandler(TransportService transportService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper,
public DefaultLwM2MUplinkMsgHandler(TransportService transportService, LwM2MAttributesService attributesService, LwM2MTransportServerConfig config, LwM2mTransportServerHelper helper,
LwM2mClientContext clientContext,
@Lazy LwM2MRpcRequestHandler rpcHandler,
@Lazy LwM2mDownlinkMsgHandler defaultLwM2MDownlinkMsgHandler,
OtaPackageDataCache otaPackageDataCache,
LwM2mTransportContext context, LwM2MJsonAdaptor adaptor, TbLwM2MDtlsSessionStore sessionStore) {
this.transportService = transportService;
this.attributesService = attributesService;
this.config = config;
this.helper = helper;
this.clientContext = clientContext;
@ -167,7 +164,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
}
@PostConstruct
public void init() {
public void initAttributes() {
this.context.getScheduler().scheduleAtFixedRate(this::reportActivity, new Random().nextInt((int) config.getSessionReportTimeout()), config.getSessionReportTimeout(), TimeUnit.MILLISECONDS);
this.registrationExecutor = ThingsBoardExecutors.newWorkStealingPool(this.config.getRegisteredPoolSize(), "LwM2M registration");
this.updateRegistrationExecutor = ThingsBoardExecutors.newWorkStealingPool(this.config.getUpdateRegisteredPoolSize(), "LwM2M update registration");
@ -202,7 +199,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
}
this.logToTelemetry(lwM2MClient, LOG_LWM2M_INFO + ": Client registered with registration id: " + registration.getId());
SessionInfoProto sessionInfo = lwM2MClient.getSession();
transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, rpcHandler, sessionInfo));
transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, attributesService, rpcHandler, sessionInfo));
log.warn("40) sessionId [{}] Registering rpc subscription after Registration client", new UUID(sessionInfo.getSessionIdMSB(), sessionInfo.getSessionIdLSB()));
TransportProtos.TransportToDeviceActorMsg msg = TransportProtos.TransportToDeviceActorMsg.newBuilder()
.setSessionInfo(sessionInfo)
@ -211,9 +208,10 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
.setSubscribeToRPC(TransportProtos.SubscribeToRPCMsg.newBuilder().setSessionType(TransportProtos.SessionType.ASYNC).build())
.build();
transportService.process(msg, null);
this.getInfoFirmwareUpdate(lwM2MClient, null);
this.getInfoSoftwareUpdate(lwM2MClient, null);
this.getInfoFirmwareUpdate(lwM2MClient);
this.getInfoSoftwareUpdate(lwM2MClient);
this.initClientTelemetry(lwM2MClient);
this.initAttributes(lwM2MClient);
} else {
log.error("Client: [{}] onRegistered [{}] name [{}] lwM2MClient ", registration.getId(), registration.getEndpoint(), null);
}
@ -324,7 +322,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
* @param response - observe
*/
@Override
public void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response, LwM2mClientRpcRequest rpcRequest) {
public void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response) {
if (response.getContent() != null) {
LwM2mClient lwM2MClient = clientContext.getClientByEndpoint(registration.getEndpoint());
ObjectModel objectModelVersion = lwM2MClient.getObjectModel(path, this.config.getModelProvider());
@ -343,70 +341,6 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
}
}
/**
* Update - send request in change value resources in Client
* 1. FirmwareUpdate:
* - If msg.getSharedUpdatedList().forEach(tsKvProto -> {tsKvProto.getKv().getKey().indexOf(FIRMWARE_UPDATE_PREFIX, 0) == 0
* 2. Shared Other AttributeUpdate
* -- Path to resources from profile equal keyName or from ModelObject equal name
* -- Only for resources: isWritable && isPresent as attribute in profile -> LwM2MClientProfile (format: CamelCase)
* 3. Delete - nothing
*
* @param msg -
*/
@Override
public void onAttributeUpdate(AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo) {
LwM2mClient lwM2MClient = clientContext.getClientBySessionInfo(sessionInfo);
if (msg.getSharedUpdatedCount() > 0 && lwM2MClient != null) {
log.warn("2) OnAttributeUpdate, SharedUpdatedList() [{}]", msg.getSharedUpdatedList());
msg.getSharedUpdatedList().forEach(tsKvProto -> {
String pathName = tsKvProto.getKv().getKey();
String pathIdVer = this.getObjectIdByKeyNameFromProfile(sessionInfo, pathName);
Object valueNew = getValueFromKvProto(tsKvProto.getKv());
if ((OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName)
&& (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion())))
|| (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.TITLE).equals(pathName)
&& (!valueNew.equals(lwM2MClient.getFwUpdate().getCurrentTitle())))) {
this.getInfoFirmwareUpdate(lwM2MClient, null);
} else if ((OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.VERSION).equals(pathName)
&& (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentVersion())))
|| (OtaPackageUtil.getAttributeKey(OtaPackageType.SOFTWARE, OtaPackageKey.TITLE).equals(pathName)
&& (!valueNew.equals(lwM2MClient.getSwUpdate().getCurrentTitle())))) {
this.getInfoSoftwareUpdate(lwM2MClient, null);
}
if (pathIdVer != null) {
ResourceModel resourceModel = lwM2MClient.getResourceModel(pathIdVer, this.config
.getModelProvider());
if (resourceModel != null && resourceModel.operations.isWritable()) {
this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer), valueNew, pathIdVer);
} else {
log.error("Resource path - [{}] value - [{}] is not Writable and cannot be updated", pathIdVer, valueNew);
String logMsg = String.format("%s: attributeUpdate: Resource path - %s value - %s is not Writable and cannot be updated",
LOG_LWM2M_ERROR, pathIdVer, valueNew);
this.logToTelemetry(lwM2MClient, logMsg);
}
} else if (!isFwSwWords(pathName)) {
log.error("Resource name name - [{}] value - [{}] is not present as attribute/telemetry in profile and cannot be updated", pathName, valueNew);
String logMsg = String.format("%s: attributeUpdate: attribute name - %s value - %s is not present as attribute in profile and cannot be updated",
LOG_LWM2M_ERROR, pathName, valueNew);
this.logToTelemetry(lwM2MClient, logMsg);
}
});
} else if (msg.getSharedDeletedCount() > 0 && lwM2MClient != null) {
msg.getSharedUpdatedList().forEach(tsKvProto -> {
String pathName = tsKvProto.getKv().getKey();
Object valueNew = getValueFromKvProto(tsKvProto.getKv());
if (OtaPackageUtil.getAttributeKey(OtaPackageType.FIRMWARE, OtaPackageKey.VERSION).equals(pathName) && !valueNew.equals(lwM2MClient.getFwUpdate().getCurrentVersion())) {
lwM2MClient.getFwUpdate().setCurrentVersion((String) valueNew);
}
});
log.info("[{}] delete [{}] onAttributeUpdate", msg.getSharedDeletedList(), sessionInfo);
} else if (lwM2MClient == null) {
log.error("OnAttributeUpdate, lwM2MClient is null");
}
}
/**
* @param sessionInfo -
* @param deviceProfile -
@ -929,7 +863,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
}
}
private void updateResourcesValueToClient(LwM2mClient lwM2MClient, Object valueOld, Object newValue, String versionedId) {
private void pushUpdateToClientIfNeeded(LwM2mClient lwM2MClient, Object valueOld, Object newValue, String versionedId) {
if (newValue != null && (valueOld == null || !newValue.toString().equals(valueOld.toString()))) {
TbLwM2MWriteReplaceRequest request = TbLwM2MWriteReplaceRequest.builder().versionedId(versionedId).value(newValue).timeout(this.config.getTimeout()).build();
defaultLwM2MDownlinkMsgHandler.sendWriteReplaceRequest(lwM2MClient, request, new TbLwM2MWriteResponseCallback(this, lwM2MClient, versionedId));
@ -973,25 +907,6 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
).getKey();
}
/**
* 1. FirmwareUpdate:
* - msg.getSharedUpdatedList().forEach(tsKvProto -> {tsKvProto.getKv().getKey().indexOf(FIRMWARE_UPDATE_PREFIX, 0) == 0
* 2. Update resource value on client: if there is a difference in values between the current resource values and the shared attribute values
* - Get path resource by result attributesResponse
*
* @param attributesResponse -
* @param sessionInfo -
*/
@Override
public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo) {
try {
List<TransportProtos.TsKvProto> tsKvProtos = attributesResponse.getSharedAttributeListList();
this.updateAttributeFromThingsboard(tsKvProtos, sessionInfo);
} catch (Exception e) {
log.error("", e);
}
}
/**
* #1.1 If two names have equal path => last time attribute
* #2.1 if there is a difference in values between the current resource values and the shared attribute values
@ -1000,31 +915,27 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
* #2.1 if there is not a difference in values between the current resource values and the shared attribute values
*
* @param tsKvProtos
* @param sessionInfo
*/
public void updateAttributeFromThingsboard(List<TransportProtos.TsKvProto> tsKvProtos, TransportProtos.SessionInfoProto sessionInfo) {
LwM2mClient lwM2MClient = clientContext.getClientBySessionInfo(sessionInfo);
if (lwM2MClient != null) {
log.warn("1) UpdateAttributeFromThingsboard, tsKvProtos [{}]", tsKvProtos);
tsKvProtos.forEach(tsKvProto -> {
String pathIdVer = this.getObjectIdByKeyNameFromProfile(sessionInfo, tsKvProto.getKv().getKey());
if (pathIdVer != null) {
// #1.1
if (lwM2MClient.getDelayedRequests().containsKey(pathIdVer) && tsKvProto.getTs() > lwM2MClient.getDelayedRequests().get(pathIdVer).getTs()) {
lwM2MClient.getDelayedRequests().put(pathIdVer, tsKvProto);
} else if (!lwM2MClient.getDelayedRequests().containsKey(pathIdVer)) {
lwM2MClient.getDelayedRequests().put(pathIdVer, tsKvProto);
public void onAttributesUpdate(LwM2mClient lwM2MClient, List<TransportProtos.TsKvProto> tsKvProtos) {
log.trace("[{}] onAttributesUpdate [{}]", lwM2MClient.getEndpoint(), tsKvProtos);
tsKvProtos.forEach(tsKvProto -> {
String pathIdVer = this.getObjectIdByKeyNameFromProfile(lwM2MClient, tsKvProto.getKv().getKey());
if (pathIdVer != null) {
// #1.1
if (lwM2MClient.getSharedAttributes().containsKey(pathIdVer)) {
if (tsKvProto.getTs() > lwM2MClient.getSharedAttributes().get(pathIdVer).getTs()) {
lwM2MClient.getSharedAttributes().put(pathIdVer, tsKvProto);
}
} else {
lwM2MClient.getSharedAttributes().put(pathIdVer, tsKvProto);
}
});
// #2.1
lwM2MClient.getDelayedRequests().forEach((pathIdVer, tsKvProto) -> {
this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer),
getValueFromKvProto(tsKvProto.getKv()), pathIdVer);
});
} else {
log.error("UpdateAttributeFromThingsboard, lwM2MClient is null");
}
}
});
// #2.1
lwM2MClient.getSharedAttributes().forEach((pathIdVer, tsKvProto) -> {
this.pushUpdateToClientIfNeeded(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer),
getValueFromKvProto(tsKvProto.getKv()), pathIdVer);
});
}
/**
@ -1053,7 +964,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
*/
private void reportActivityAndRegister(SessionInfoProto sessionInfo) {
if (sessionInfo != null && transportService.reportActivity(sessionInfo) == null) {
transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, rpcHandler, sessionInfo));
transportService.registerAsyncSession(sessionInfo, new LwM2mSessionMsgListener(this, attributesService, rpcHandler, sessionInfo));
this.reportActivitySubscription(sessionInfo);
}
}
@ -1072,25 +983,20 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
*
* @param lwM2MClient - LwM2M Client
*/
public void putDelayedUpdateResourcesThingsboard(LwM2mClient lwM2MClient) {
SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient);
if (sessionInfo != null) {
//#1.1
Map<String, String> keyNamesMap = this.getNamesFromProfileForSharedAttributes(lwM2MClient);
if (keyNamesMap.values().size() > 0) {
try {
//#1.2
TransportProtos.GetAttributeRequestMsg getAttributeMsg = adaptor.convertToGetAttributes(null, keyNamesMap.values());
transportService.process(sessionInfo, getAttributeMsg, getAckCallback(lwM2MClient, getAttributeMsg.getRequestId(), DEVICE_ATTRIBUTES_REQUEST));
} catch (AdaptorException e) {
log.trace("Failed to decode get attributes request", e);
}
}
public void initAttributes(LwM2mClient lwM2MClient) {
Map<String, String> keyNamesMap = this.getNamesFromProfileForSharedAttributes(lwM2MClient);
if (!keyNamesMap.isEmpty()) {
Set<String> keysToFetch = new HashSet<>(keyNamesMap.values());
keysToFetch.removeAll(OtaPackageUtil.ALL_FW_ATTRIBUTE_KEYS);
keysToFetch.removeAll(OtaPackageUtil.ALL_SW_ATTRIBUTE_KEYS);
DonAsynchron.withCallback(attributesService.getSharedAttributes(lwM2MClient, keysToFetch),
v -> onAttributesUpdate(lwM2MClient, v),
t -> log.error("[{}] Failed to get attributes", lwM2MClient.getEndpoint(), t),
registrationExecutor);
}
}
public void getInfoFirmwareUpdate(LwM2mClient lwM2MClient, LwM2mClientRpcRequest rpcRequest) {
public void getInfoFirmwareUpdate(LwM2mClient lwM2MClient) {
if (lwM2MClient.getRegistration().getSupportedVersion(FW_5_ID) != null) {
SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient);
if (sessionInfo != null) {
@ -1102,20 +1008,20 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
if (TransportProtos.ResponseStatus.SUCCESS.equals(response.getResponseStatus())
&& response.getType().equals(OtaPackageType.FIRMWARE.name())) {
LwM2mFwSwUpdate fwUpdate = lwM2MClient.getFwUpdate(DefaultLwM2MUplinkMsgHandler.this, clientContext);
if (rpcRequest != null) {
fwUpdate.setStateUpdate(INITIATED.name());
}
// if (rpcRequest != null) {
// fwUpdate.setStateUpdate(INITIATED.name());
// }
if (!FAILED.name().equals(fwUpdate.getStateUpdate())) {
log.warn("7) firmware start with ver: [{}]", response.getVersion());
fwUpdate.setRpcRequest(rpcRequest);
// fwUpdate.setRpcRequest(rpcRequest);
fwUpdate.setCurrentVersion(response.getVersion());
fwUpdate.setCurrentTitle(response.getTitle());
fwUpdate.setCurrentId(new UUID(response.getOtaPackageIdMSB(), response.getOtaPackageIdLSB()));
if (rpcRequest == null) {
// if (rpcRequest == null) {
fwUpdate.sendReadObserveInfo(defaultLwM2MDownlinkMsgHandler);
} else {
fwUpdate.writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler);
}
// } else {
// fwUpdate.writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler);
// }
} else {
String msgError = String.format("OtaPackage device: %s, version: %s, stateUpdate: %s",
lwM2MClient.getDeviceName(), response.getVersion(), fwUpdate.getStateUpdate());
@ -1125,10 +1031,10 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
String msgError = String.format("OtaPackage device: %s, responseStatus: %s",
lwM2MClient.getDeviceName(), response.getResponseStatus().toString());
log.trace(msgError);
if (rpcRequest != null) {
// if (rpcRequest != null) {
//TODO: refactor
// sendErrorRpcResponse(rpcRequest, msgError, sessionInfo);
}
// }
}
}
@ -1141,7 +1047,7 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
}
}
public void getInfoSoftwareUpdate(LwM2mClient lwM2MClient, LwM2mClientRpcRequest rpcRequest) {
public void getInfoSoftwareUpdate(LwM2mClient lwM2MClient) {
if (lwM2MClient.getRegistration().getSupportedVersion(SW_ID) != null) {
SessionInfoProto sessionInfo = this.getSessionInfo(lwM2MClient);
if (sessionInfo != null) {
@ -1152,16 +1058,16 @@ public class DefaultLwM2MUplinkMsgHandler implements LwM2mUplinkMsgHandler {
public void onSuccess(TransportProtos.GetOtaPackageResponseMsg response) {
if (TransportProtos.ResponseStatus.SUCCESS.equals(response.getResponseStatus())
&& response.getType().equals(OtaPackageType.SOFTWARE.name())) {
lwM2MClient.getSwUpdate().setRpcRequest(rpcRequest);
// lwM2MClient.getSwUpdate().setRpcRequest(rpcRequest);
lwM2MClient.getSwUpdate().setCurrentVersion(response.getVersion());
lwM2MClient.getSwUpdate().setCurrentTitle(response.getTitle());
lwM2MClient.getSwUpdate().setCurrentId(new OtaPackageId(new UUID(response.getOtaPackageIdMSB(), response.getOtaPackageIdLSB())).getId());
lwM2MClient.getSwUpdate().sendReadObserveInfo(defaultLwM2MDownlinkMsgHandler);
if (rpcRequest == null) {
// if (rpcRequest == null) {
lwM2MClient.getSwUpdate().sendReadObserveInfo(defaultLwM2MDownlinkMsgHandler);
} else {
lwM2MClient.getSwUpdate().writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler);
}
// } else {
// lwM2MClient.getSwUpdate().writeFwSwWare(handler, defaultLwM2MDownlinkMsgHandler);
// }
} else {
log.trace("Software [{}] [{}]", lwM2MClient.getDeviceName(), response.getResponseStatus().toString());
}

9
common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/uplink/LwM2mUplinkMsgHandler.java

@ -24,7 +24,6 @@ import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.transport.lwm2m.config.LwM2MTransportServerConfig;
import org.thingsboard.server.transport.lwm2m.server.client.LwM2mClient;
import org.thingsboard.server.transport.lwm2m.server.rpc.LwM2mClientRpcRequest;
import java.util.Collection;
import java.util.Optional;
@ -39,15 +38,13 @@ public interface LwM2mUplinkMsgHandler {
void onSleepingDev(Registration registration);
void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response, LwM2mClientRpcRequest rpcRequest);
void onAttributeUpdate(TransportProtos.AttributeUpdateNotificationMsg msg, TransportProtos.SessionInfoProto sessionInfo);
void onUpdateValueAfterReadResponse(Registration registration, String path, ReadResponse response);
void onDeviceProfileUpdate(TransportProtos.SessionInfoProto sessionInfo, DeviceProfile deviceProfile);
void onDeviceUpdate(TransportProtos.SessionInfoProto sessionInfo, Device device, Optional<DeviceProfile> deviceProfileOpt);
void onResourceUpdate (Optional<TransportProtos.ResourceUpdateMsg> resourceUpdateMsgOpt);
void onResourceUpdate(Optional<TransportProtos.ResourceUpdateMsg> resourceUpdateMsgOpt);
void onResourceDelete(Optional<TransportProtos.ResourceDeleteMsg> resourceDeleteMsgOpt);
@ -67,7 +64,5 @@ public interface LwM2mUplinkMsgHandler {
String getObjectIdByKeyNameFromProfile(LwM2mClient lwM2mClient, String keyName);
void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo);
LwM2MTransportServerConfig getConfig();
}

Loading…
Cancel
Save