resourceUpdateMsgOpt) {
String idVer = resourceUpdateMsgOpt.get().getResourceKey();
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.updateResourceModel(idVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider()));
}
/**
- *
* @param resourceDeleteMsgOpt -
*/
@Override
@@ -375,6 +381,14 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.deleteResources(pathIdVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider()));
}
+ public void onToDeviceRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRequest) {
+ log.info("[{}] toDeviceRpcRequest", toDeviceRequest);
+ }
+
+ public void onToServerRpcResponse(TransportProtos.ToServerRpcResponseMsg toServerResponse) {
+ log.info("[{}] toServerRpcResponse", toServerResponse);
+ }
+
/**
* Trigger Server path = "/1/0/8"
*
@@ -418,6 +432,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
*/
protected void onAwakeDev(Registration registration) {
log.info("[{}] [{}] Received endpoint Awake version event", registration.getId(), registration.getEndpoint());
+ this.sendLogsToThingsboard(LOG_LW2M_INFO + ": Client is awake!", registration);
//TODO: associate endpointId with device information.
}
@@ -443,32 +458,13 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
}
/**
- * @param msg - text msg
+ * @param logMsg - text msg
* @param registration - Id of Registration LwM2M Client
*/
- public void sendLogsToThingsboard(String msg, Registration registration) {
- if (msg != null) {
- JsonObject telemetries = new JsonObject();
- telemetries.addProperty(LOG_LW2M_TELEMETRY, msg);
- this.updateParametersOnThingsboard(telemetries, DEVICE_TELEMETRY_TOPIC, registration);
- }
- }
-
-
- /**
- * // !!! Ok
- * Prepare send to Thigsboard callback - Attribute or Telemetry
- *
- * @param msg - JsonArray: [{name: value}]
- * @param topicName - Api Attribute or Telemetry
- * @param registration - Id of Registration LwM2M Client
- */
- public void updateParametersOnThingsboard(JsonElement msg, String topicName, Registration registration) {
+ public void sendLogsToThingsboard(String logMsg, Registration registration) {
SessionInfoProto sessionInfo = this.getValidateSessionInfo(registration);
- if (sessionInfo != null) {
- lwM2mTransportContextServer.sendParametersOnThingsboard(msg, topicName, sessionInfo);
- } else {
- log.error("Client: [{}] updateParametersOnThingsboard [{}] sessionInfo ", registration, null);
+ if (logMsg != null && sessionInfo != null) {
+ this.lwM2mTransportContextServer.sendParametersOnThingsboardTelemetry(this.lwM2mTransportContextServer.getKvLogyToThingsboard(logMsg), sessionInfo);
}
}
@@ -486,15 +482,20 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
*/
private void initLwM2mFromClientValue(Registration registration, LwM2mClient lwM2MClient) {
LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration);
- Set clientObjects = this.getAllOjectsInClient(registration);
- if (clientObjects != null && LWM2M_STRATEGY_2 == LwM2mTransportHandler.getClientOnlyObserveAfterConnect(lwM2MClientProfile)) {
- // #2
- lwM2MClient.getPendingRequests().addAll(clientObjects);
- clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(registration, path, GET_TYPE_OPER_READ, ContentFormat.TLV.getName(),
- null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
+ Set clientObjects = lwM2mClientContext.getSupportedIdVerInClient(registration);
+ if (clientObjects != null && clientObjects.size() > 0) {
+ if (LWM2M_STRATEGY_2 == LwM2mTransportHandler.getClientOnlyObserveAfterConnect(lwM2MClientProfile)) {
+ // #2
+ lwM2MClient.getPendingRequests().addAll(clientObjects);
+ clientObjects.forEach(path -> lwM2mTransportRequest.sendAllRequest(registration, path, GET_TYPE_OPER_READ, ContentFormat.TLV.getName(),
+ null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
+ }
+ // #1
+ this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_READ, clientObjects);
+ this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_OBSERVE, clientObjects);
+ this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, PUT_TYPE_OPER_WRITE_ATTRIBUTES, clientObjects);
+ this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_DISCOVER, clientObjects);
}
- // #1
- this.initReadAttrTelemetryObserveToClient(registration, lwM2MClient, GET_TYPE_OPER_OBSERVE);
}
/**
@@ -553,17 +554,20 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* @param registration - Registration LwM2M Client
*/
private void updateAttrTelemetry(Registration registration, Set paths) {
- JsonObject attributes = new JsonObject();
- JsonObject telemetries = new JsonObject();
try {
- this.getParametersFromProfile(attributes, telemetries, registration, paths);
+ ResultsAddKeyValueProto results = getParametersFromProfile(registration, paths);
+ SessionInfoProto sessionInfo = this.getValidateSessionInfo(registration);
+ if (results != null && sessionInfo != null) {
+ if (results.getResultAttributes().size() > 0) {
+ this.lwM2mTransportContextServer.sendParametersOnThingsboardAttribute(results.getResultAttributes(), sessionInfo);
+ }
+ if (results.getResultTelemetries().size() > 0) {
+ this.lwM2mTransportContextServer.sendParametersOnThingsboardTelemetry(results.getResultTelemetries(), sessionInfo);
+ }
+ }
} catch (Exception e) {
log.error("UpdateAttrTelemetry", e);
}
- if (attributes.getAsJsonObject().entrySet().size() > 0)
- this.updateParametersOnThingsboard(attributes, DEVICE_ATTRIBUTES_TOPIC, registration);
- if (telemetries.getAsJsonObject().entrySet().size() > 0)
- this.updateParametersOnThingsboard(telemetries, DEVICE_TELEMETRY_TOPIC, registration);
}
/**
@@ -601,41 +605,61 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
/**
* Start observe/read: Attr/Telemetry
- * #1 - Analyze:
- * #1.1 path in resource profile == client resource
+ * #1 - Analyze: path in resource profile == client resource
*
* @param registration -
*/
- private void initReadAttrTelemetryObserveToClient(Registration registration, LwM2mClient lwM2MClient, String typeOper) {
+ private void initReadAttrTelemetryObserveToClient(Registration registration, LwM2mClient lwM2MClient,
+ String typeOper, Set clientObjects) {
LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration);
- Set clientInstances = this.getAllInstancesInClient(registration);
- Set result;
+ Set result = null;
+ ConcurrentHashMap params = null;
if (GET_TYPE_OPER_READ.equals(typeOper)) {
- result = JacksonUtil.fromString(lwM2MClientProfile.getPostAttributeProfile().toString(), new TypeReference<>() {
- });
- result.addAll(JacksonUtil.fromString(lwM2MClientProfile.getPostTelemetryProfile().toString(), new TypeReference<>() {
- }));
- } else {
- result = JacksonUtil.fromString(lwM2MClientProfile.getPostObserveProfile().toString(), new TypeReference<>() {
- });
+ result = JacksonUtil.fromString(lwM2MClientProfile.getPostAttributeProfile().toString(),
+ new TypeReference<>() {
+ });
+ result.addAll(JacksonUtil.fromString(lwM2MClientProfile.getPostTelemetryProfile().toString(),
+ new TypeReference<>() {
+ }));
+ } else if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) {
+ result = JacksonUtil.fromString(lwM2MClientProfile.getPostObserveProfile().toString(),
+ new TypeReference<>() {
+ });
+ } else if (GET_TYPE_OPER_DISCOVER.equals(typeOper)) {
+ result = this.getPathForWriteAttributes(lwM2MClientProfile.getPostAttributeLwm2mProfile()).keySet();
+ ;
+ } else if (PUT_TYPE_OPER_WRITE_ATTRIBUTES.equals(typeOper)) {
+ params = this.getPathForWriteAttributes(lwM2MClientProfile.getPostAttributeLwm2mProfile());
+ result = params.keySet();
}
- Set pathSend = ConcurrentHashMap.newKeySet();
- result.forEach(target -> {
- // #1.1
- String[] resPath = target.split("/");
- String instance = "/" + resPath[1] + "/" + resPath[2];
- if (clientInstances != null && clientInstances.size() > 0 && clientInstances.contains(instance)) {
- pathSend.add(target);
+ if (!result.isEmpty()) {
+ // #1
+ Set pathSend = result.stream().filter(target -> {
+ return target.split(LWM2M_SEPARATOR_PATH).length < 3 ?
+ clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1]) :
+ clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1] + "/" + target.split(LWM2M_SEPARATOR_PATH)[2]);
+ }
+ )
+ .collect(Collectors.toUnmodifiableSet());
+ if (!pathSend.isEmpty()) {
+ lwM2MClient.getPendingRequests().addAll(pathSend);
+ ConcurrentHashMap finalParams = params;
+ pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, ContentFormat.TLV.getName(),
+ null, finalParams != null ? finalParams.get(target) : null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
+ if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) {
+ lwM2MClient.initValue(this, null);
+ }
}
- });
- lwM2MClient.getPendingRequests().addAll(pathSend);
- pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, typeOper, ContentFormat.TLV.getName(),
- null, null, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
- if (GET_TYPE_OPER_OBSERVE.equals(typeOper)) {
- lwM2MClient.initValue(this, null);
}
}
+ private ConcurrentHashMap getPathForWriteAttributes(JsonObject objectJson) {
+ ConcurrentHashMap pathAttributes = new Gson().fromJson(objectJson.toString(),
+ new TypeToken>() {
+ }.getType());
+ return pathAttributes;
+ }
+
/**
* Update parameters device in LwM2MClient
* If new deviceProfile != old deviceProfile => update deviceProfile
@@ -656,96 +680,109 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
}
/**
- * @param registration -
- * @return - all object in client
- */
- private Set getAllOjectsInClient(Registration registration) {
- Set clientObjects = ConcurrentHashMap.newKeySet();
- Arrays.stream(registration.getObjectLinks()).forEach(url -> {
- LwM2mPath pathIds = new LwM2mPath(url.getUrl());
- if (pathIds.isObjectInstance()) {
- clientObjects.add("/" + pathIds.getObjectId());
- }
- });
- return (clientObjects.size() > 0) ? clientObjects : null;
- }
-
- /**
- * @param registration -
- * @return all instances in client
- */
- private Set getAllInstancesInClient(Registration registration) {
- Set clientInstances = ConcurrentHashMap.newKeySet();
- Arrays.stream(registration.getObjectLinks()).forEach(url -> {
- LwM2mPath pathIds = new LwM2mPath(url.getUrl());
- if (pathIds.isObjectInstance()) {
- clientInstances.add(convertToIdVerFromObjectId(url.getUrl(), registration));
- }
- });
- return (clientInstances.size() > 0) ? clientInstances : null;
- }
-
- /**
- * @param attributes - new JsonObject
- * @param telemetry - new JsonObject
+ * // * @param attributes - new JsonObject
+ * // * @param telemetry - new JsonObject
+ *
* @param registration - Registration LwM2M Client
* @param path -
*/
- private void getParametersFromProfile(JsonObject attributes, JsonObject telemetry, Registration registration, Set path) {
+ private ResultsAddKeyValueProto getParametersFromProfile(Registration registration, Set path) {
if (path != null && path.size() > 0) {
+ ResultsAddKeyValueProto results = new ResultsAddKeyValueProto();
LwM2mClientProfile lwM2MClientProfile = lwM2mClientContext.getProfile(registration);
- lwM2MClientProfile.getPostAttributeProfile().forEach(idVer -> {
- if (path.contains(idVer.getAsString())) {
- this.addParameters(idVer.getAsString(), attributes, registration);
+ List resultAttributes = new ArrayList<>();
+ lwM2MClientProfile.getPostAttributeProfile().forEach(pathIdVer -> {
+ if (path.contains(pathIdVer.getAsString())) {
+ TransportProtos.KeyValueProto kvAttr = this.getKvToThingsboard(pathIdVer.getAsString(), registration);
+ if (kvAttr != null) {
+ resultAttributes.add(kvAttr);
+ }
}
});
- lwM2MClientProfile.getPostTelemetryProfile().forEach(idVer -> {
- if (path.contains(idVer.getAsString())) {
- this.addParameters(idVer.getAsString(), telemetry, registration);
+ List resultTelemetries = new ArrayList<>();
+ lwM2MClientProfile.getPostTelemetryProfile().forEach(pathIdVer -> {
+ if (path.contains(pathIdVer.getAsString())) {
+ TransportProtos.KeyValueProto kvAttr = this.getKvToThingsboard(pathIdVer.getAsString(), registration);
+ if (kvAttr != null) {
+ resultTelemetries.add(kvAttr);
+ }
}
});
+ if (resultAttributes.size() > 0) {
+ results.setResultAttributes(resultAttributes);
+ }
+ if (resultTelemetries.size() > 0) {
+ results.setResultTelemetries(resultTelemetries);
+ }
+ return results;
}
+ return null;
}
- /**
- * @param parameters - JsonObject attributes/telemetry
- * @param registration - Registration LwM2M Client
- */
- private void addParameters(String path, JsonObject parameters, Registration registration) {
- LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null);
+// public TransportProtos.KeyValueProto getKvToThingsboard(String pathIdVer, Registration registration) {
+// ResultsResourceValue resultsResourceValue = getResultsResourceValue(pathIdVer, registration);
+// return resultsResourceValue != null ? this.lwM2mTransportContextServer.getKvAttrTelemetryToThingsboard(resultsResourceValue) : null;
+// }
+
+ private TransportProtos.KeyValueProto getKvToThingsboard(String pathIdVer, Registration registration) {
+ LwM2mClient lwM2MClient = this.lwM2mClientContext.getLwM2mClientWithReg(null, registration.getId());
JsonObject names = lwM2mClientContext.getProfiles().get(lwM2MClient.getProfileId()).getPostKeyNameProfile();
- String resName = names.get(path).getAsString();
- if (resName != null && !resName.isEmpty()) {
- try {
- String resValue = this.getResourceValueToString(lwM2MClient, path);
- if (resValue != null) {
- parameters.addProperty(resName, resValue);
+ if (names != null && names.has(pathIdVer)) {
+ String resourceName = names.get(pathIdVer).getAsString();
+ if (resourceName != null && !resourceName.isEmpty()) {
+ try {
+ LwM2mResource resourceValue = lwM2MClient != null ? getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer) : null;
+ if (resourceValue != null) {
+ ResourceModel.Type currentType = resourceValue.getType();
+ ResourceModel.Type expectedType = this.lwM2mTransportContextServer.getResourceModelTypeEqualsKvProtoValueType(currentType, pathIdVer);
+ Object valueKvProto = null;
+ if (resourceValue.isMultiInstances()) {
+ valueKvProto = new JsonObject();
+ Object finalvalueKvProto = valueKvProto;
+ Gson gson = new GsonBuilder().create();
+ resourceValue.getValues().forEach((k, v) -> {
+ Object val = this.converter.convertValue(resourceValue.getValue(), currentType, expectedType,
+ new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)));
+ JsonElement element = gson.toJsonTree(val, val.getClass());
+ ((JsonObject) finalvalueKvProto).add(String.valueOf(k), element);
+ });
+ valueKvProto = gson.toJson(valueKvProto);
+ } else {
+ valueKvProto = this.converter.convertValue(resourceValue.getValue(), currentType, expectedType,
+ new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)));
+ }
+ return valueKvProto != null ? this.lwM2mTransportContextServer.getKvAttrTelemetryToThingsboard(currentType, resourceName, valueKvProto, resourceValue.isMultiInstances()) : null;
+ }
+ } catch (Exception e) {
+ log.error("Failed to add parameters.", e);
}
- } catch (Exception e) {
- log.error("Failed to add parameters.", e);
}
+ } else {
+ log.error("Failed to add parameters. path: [{}], names: [{}]", pathIdVer, names);
}
+ return null;
}
/**
- * @param path - path resource
- * @return - value of Resource or null
+ * @param pathIdVer - path resource
+ * @return - value of Resource into format KvProto or null
*/
- private String getResourceValueToString(LwM2mClient lwM2MClient, String path) {
- LwM2mPath pathIds = new LwM2mPath(convertToObjectIdFromIdVer(path));
- LwM2mResource resourceValue = this.returnResourceValueFromLwM2MClient(lwM2MClient, path);
- return resourceValue == null ? null :
- this.converter.convertValue(resourceValue.isMultiInstances() ? resourceValue.getValues() : resourceValue.getValue(), resourceValue.getType(), ResourceModel.Type.STRING, pathIds).toString();
+ private Object getResourceValueFormatKv(LwM2mClient lwM2MClient, String pathIdVer) {
+ LwM2mResource resourceValue = this.getResourceValueFromLwM2MClient(lwM2MClient, pathIdVer);
+ ResourceModel.Type currentType = resourceValue.getType();
+ ResourceModel.Type expectedType = this.lwM2mTransportContextServer.getResourceModelTypeEqualsKvProtoValueType(currentType, pathIdVer);
+ return this.converter.convertValue(resourceValue.getValue(), currentType, expectedType,
+ new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)));
}
/**
* @param lwM2MClient -
- * @param path -
+ * @param path -
* @return - return value of Resource by idPath
*/
- private LwM2mResource returnResourceValueFromLwM2MClient(LwM2mClient lwM2MClient, String path) {
+ private LwM2mResource getResourceValueFromLwM2MClient(LwM2mClient lwM2MClient, String path) {
LwM2mResource resourceValue = null;
- if (new LwM2mPath(convertToObjectIdFromIdVer(path)).isResource()) {
+ if (new LwM2mPath(convertPathFromIdVerToObjectId(path)).isResource()) {
resourceValue = lwM2MClient.getResources().get(path).getLwM2mResource();
}
return resourceValue;
@@ -769,6 +806,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* #3.1 Attribute isChange (add&del)
* #3.2 Telemetry isChange (add&del)
* #3.3 KeyName isChange (add)
+ * #3.4 attributeLwm2m isChange (update WrightAttribute: add/update/del)
* #4 update
* #4.1 add If #3 isChange, then analyze and update Value in Transport form Client and send Value to thingsboard
* #4.2 del
@@ -779,6 +817,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* -- path Attr/Telemetry includes newObserve and does not include oldObserve: send Request observe to Client
* #5.3 Observe.del
* -- different between newObserve and oldObserve: send Request cancel observe to client
+ * #6
+ * #6.1 - update WriteAttribute
+ * #6.2 - del WriteAttribute
*
* @param registrationIds -
* @param deviceProfile -
@@ -793,6 +834,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
Set telemetrySetOld = this.convertJsonArrayToSet(telemetryOld);
JsonArray observeOld = lwM2MClientProfileOld.getPostObserveProfile();
JsonObject keyNameOld = lwM2MClientProfileOld.getPostKeyNameProfile();
+ JsonObject attributeLwm2mOld = lwM2MClientProfileOld.getPostAttributeLwm2mProfile();
LwM2mClientProfile lwM2MClientProfileNew = lwM2mClientContext.getProfiles().get(deviceProfile.getUuidId());
JsonArray attributeNew = lwM2MClientProfileNew.getPostAttributeProfile();
@@ -801,32 +843,41 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
Set telemetrySetNew = this.convertJsonArrayToSet(telemetryNew);
JsonArray observeNew = lwM2MClientProfileNew.getPostObserveProfile();
JsonObject keyNameNew = lwM2MClientProfileNew.getPostKeyNameProfile();
+ JsonObject attributeLwm2mNew = lwM2MClientProfileNew.getPostAttributeLwm2mProfile();
// #3
ResultsAnalyzerParameters sendAttrToThingsboard = new ResultsAnalyzerParameters();
// #3.1
if (!attributeOld.equals(attributeNew)) {
- ResultsAnalyzerParameters postAttributeAnalyzer = this.getAnalyzerParameters(new Gson().fromJson(attributeOld, new TypeToken>() {
- }.getType()), attributeSetNew);
+ ResultsAnalyzerParameters postAttributeAnalyzer = this.getAnalyzerParameters(new Gson().fromJson(attributeOld,
+ new TypeToken>() {
+ }.getType()), attributeSetNew);
sendAttrToThingsboard.getPathPostParametersAdd().addAll(postAttributeAnalyzer.getPathPostParametersAdd());
sendAttrToThingsboard.getPathPostParametersDel().addAll(postAttributeAnalyzer.getPathPostParametersDel());
}
// #3.2
if (!telemetryOld.equals(telemetryNew)) {
- ResultsAnalyzerParameters postTelemetryAnalyzer = this.getAnalyzerParameters(new Gson().fromJson(telemetryOld, new TypeToken>() {
- }.getType()), telemetrySetNew);
+ ResultsAnalyzerParameters postTelemetryAnalyzer = this.getAnalyzerParameters(new Gson().fromJson(telemetryOld,
+ new TypeToken>() {
+ }.getType()), telemetrySetNew);
sendAttrToThingsboard.getPathPostParametersAdd().addAll(postTelemetryAnalyzer.getPathPostParametersAdd());
sendAttrToThingsboard.getPathPostParametersDel().addAll(postTelemetryAnalyzer.getPathPostParametersDel());
}
// #3.3
if (!keyNameOld.equals(keyNameNew)) {
- ResultsAnalyzerParameters keyNameChange = this.getAnalyzerKeyName(new Gson().fromJson(keyNameOld.toString(), new TypeToken>() {
+ ResultsAnalyzerParameters keyNameChange = this.getAnalyzerKeyName(new Gson().fromJson(keyNameOld.toString(),
+ new TypeToken>() {
}.getType()),
new Gson().fromJson(keyNameNew.toString(), new TypeToken>() {
}.getType()));
sendAttrToThingsboard.getPathPostParametersAdd().addAll(keyNameChange.getPathPostParametersAdd());
}
+ // #3.4, #6
+ if (!attributeLwm2mOld.equals(attributeLwm2mNew)) {
+ this.getAnalyzerAttributeLwm2m(registrationIds, attributeLwm2mOld, attributeLwm2mNew);
+ }
+
// #4.1 add
if (sendAttrToThingsboard.getPathPostParametersAdd().size() > 0) {
// update value in Resources
@@ -860,10 +911,14 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
// send Request observe to Client
registrationIds.forEach(registrationId -> {
Registration registration = lwM2mClientContext.getRegistration(registrationId);
- this.readResourceValueObserve(registration, postObserveAnalyzer.getPathPostParametersAdd(), GET_TYPE_OPER_OBSERVE);
+ if (postObserveAnalyzer.getPathPostParametersAdd().size() > 0) {
+ this.readResourceValueObserve(registration, postObserveAnalyzer.getPathPostParametersAdd(), GET_TYPE_OPER_OBSERVE);
+ }
// 5.3 del
// send Request cancel observe to Client
- this.cancelObserveIsValue(registration, postObserveAnalyzer.getPathPostParametersDel());
+ if (postObserveAnalyzer.getPathPostParametersDel().size() > 0) {
+ this.cancelObserveIsValue(registration, postObserveAnalyzer.getPathPostParametersDel());
+ }
});
}
}
@@ -910,7 +965,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
*/
private void readResourceValueObserve(Registration registration, Set targets, String typeOper) {
targets.forEach(target -> {
- LwM2mPath pathIds = new LwM2mPath(convertToObjectIdFromIdVer(target));
+ LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(target));
if (pathIds.isResource()) {
if (GET_TYPE_OPER_READ.equals(typeOper)) {
lwM2mTransportRequest.sendAllRequest(registration, target, typeOper,
@@ -923,7 +978,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
});
}
- private ResultsAnalyzerParameters getAnalyzerKeyName(ConcurrentMap keyNameOld, ConcurrentMap keyNameNew) {
+ private ResultsAnalyzerParameters getAnalyzerKeyName(ConcurrentHashMap keyNameOld, ConcurrentHashMap keyNameNew) {
ResultsAnalyzerParameters analyzerParameters = new ResultsAnalyzerParameters();
Set paths = keyNameNew.entrySet()
.stream()
@@ -933,22 +988,91 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
return analyzerParameters;
}
+ /**
+ * #3.4, #6
+ * #6
+ * #6.1 - send update WriteAttribute
+ * #6.2 - send empty WriteAttribute
+ *
+ * @param attributeLwm2mOld -
+ * @param attributeLwm2mNew -
+ * @return
+ */
+ private void getAnalyzerAttributeLwm2m(Set registrationIds, JsonObject attributeLwm2mOld, JsonObject attributeLwm2mNew) {
+ ResultsAnalyzerParameters analyzerParameters = new ResultsAnalyzerParameters();
+ ConcurrentHashMap lwm2mAttributesOld = new Gson().fromJson(attributeLwm2mOld.toString(),
+ new TypeToken>() {
+ }.getType());
+ ConcurrentHashMap lwm2mAttributesNew = new Gson().fromJson(attributeLwm2mNew.toString(),
+ new TypeToken>() {
+ }.getType());
+ Set pathOld = lwm2mAttributesOld.keySet();
+ Set pathNew = lwm2mAttributesNew.keySet();
+ analyzerParameters.setPathPostParametersAdd(pathNew
+ .stream().filter(p -> !pathOld.contains(p)).collect(Collectors.toSet()));
+ analyzerParameters.setPathPostParametersDel(pathOld
+ .stream().filter(p -> !pathNew.contains(p)).collect(Collectors.toSet()));
+ Set pathCommon = pathNew
+ .stream().filter(p -> pathOld.contains(p)).collect(Collectors.toSet());
+ Set pathCommonChange = pathCommon
+ .stream().filter(p -> !lwm2mAttributesOld.get(p).equals(lwm2mAttributesNew.get(p))).collect(Collectors.toSet());
+ analyzerParameters.getPathPostParametersAdd().addAll(pathCommonChange);
+ // #6
+ // #6.2
+ if (analyzerParameters.getPathPostParametersAdd().size() > 0) {
+ registrationIds.forEach(registrationId -> {
+ Registration registration = this.lwM2mClientContext.getRegistration(registrationId);
+ Set clientObjects = lwM2mClientContext.getSupportedIdVerInClient(registration);
+ Set pathSend = analyzerParameters.getPathPostParametersAdd().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1]))
+ .collect(Collectors.toUnmodifiableSet());
+ if (!pathSend.isEmpty()) {
+ ConcurrentHashMap finalParams = lwm2mAttributesNew;
+ pathSend.forEach(target -> lwM2mTransportRequest.sendAllRequest(registration, target, PUT_TYPE_OPER_WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
+ null, finalParams.get(target), this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout()));
+ }
+ });
+ }
+ // #6.2
+ if (analyzerParameters.getPathPostParametersDel().size() > 0) {
+ registrationIds.forEach(registrationId -> {
+ Registration registration = this.lwM2mClientContext.getRegistration(registrationId);
+ Set clientObjects = lwM2mClientContext.getSupportedIdVerInClient(registration);
+ Set pathSend = analyzerParameters.getPathPostParametersDel().stream().filter(target -> clientObjects.contains("/" + target.split(LWM2M_SEPARATOR_PATH)[1]))
+ .collect(Collectors.toUnmodifiableSet());
+ if (!pathSend.isEmpty()) {
+ pathSend.forEach(target -> {
+ Map params = (Map) lwm2mAttributesOld.get(target);
+ params.clear();
+ params.put(OBJECT_VERSION, "");
+ lwM2mTransportRequest.sendAllRequest(registration, target, PUT_TYPE_OPER_WRITE_ATTRIBUTES, ContentFormat.TLV.getName(),
+ null, params, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout());
+ });
+ }
+ });
+ }
+
+ }
+
private void cancelObserveIsValue(Registration registration, Set paramAnallyzer) {
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null);
paramAnallyzer.forEach(p -> {
- if (this.returnResourceValueFromLwM2MClient(lwM2MClient, p) != null) {
- this.setCancelObservationRecourse(registration, convertToObjectIdFromIdVer(p));
+ if (this.getResourceValueFromLwM2MClient(lwM2MClient, p) != null) {
+ this.setCancelObservationRecourse(registration, convertPathFromIdVerToObjectId(p));
}
}
);
}
- private void putDelayedUpdateResourcesClient(LwM2mClient lwM2MClient, Object valueOld, Object valueNew, String path) {
+ private void updateResourcesValueToClient(LwM2mClient lwM2MClient, Object valueOld, Object valueNew, String path) {
if (valueNew != null && (valueOld == null || !valueNew.toString().equals(valueOld.toString()))) {
lwM2mTransportRequest.sendAllRequest(lwM2MClient.getRegistration(), path, POST_TYPE_OPER_WRITE_REPLACE,
ContentFormat.TLV.getName(), null, valueNew, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getTimeout());
} else {
- log.error("05 delayError");
+ log.error("Failed update resource [{}] [{}]", path, valueNew);
+ String logMsg = String.format("%s: Failed update resource path - %s value - %s. Value is not changed or bad",
+ LOG_LW2M_ERROR, path, valueNew);
+ this.sendLogsToThingsboard(logMsg, lwM2MClient.getRegistration());
+ log.info("Failed update resource [{}] [{}]", path, valueNew);
}
}
@@ -968,9 +1092,9 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* @param name -
* @return path if path isPresent in postProfile
*/
- private String getPathAttributeUpdate(TransportProtos.SessionInfoProto sessionInfo, String name) {
- String profilePath = this.getPathAttributeUpdateProfile(sessionInfo, name);
- return !profilePath.isEmpty() ? profilePath : null;
+ private String validatePathIntoProfile(TransportProtos.SessionInfoProto sessionInfo, String name) {
+ String pathIdVer = this.getPresentPathIntoProfile(sessionInfo, name);
+ return !pathIdVer.isEmpty() ? pathIdVer : null;
}
/**
@@ -980,7 +1104,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
* @param name -
* @return -
*/
- private String getPathAttributeUpdateProfile(TransportProtos.SessionInfoProto sessionInfo, String name) {
+ private String getPresentPathIntoProfile(TransportProtos.SessionInfoProto sessionInfo, String name) {
LwM2mClientProfile profile = lwM2mClientContext.getProfile(new UUID(sessionInfo.getDeviceProfileIdMSB(), sessionInfo.getDeviceProfileIdLSB()));
LwM2mClient lwM2mClient = lwM2mClientContext.getLwM2MClient(sessionInfo);
return profile.getPostKeyNameProfile().getAsJsonObject().entrySet().stream()
@@ -991,40 +1115,49 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
/**
* Update resource value on client: if there is a difference in values between the current resource values and the shared attribute values
* #1 Get path resource by result attributesResponse
- * #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
- * => send to client Request Update of value (new value from shared attribute)
- * and LwM2MClient.delayedRequests.add(path)
- * #2.1 if there is not a difference in values between the current resource values and the shared attribute values
*
* @param attributesResponse -
* @param sessionInfo -
*/
public void onGetAttributesResponse(TransportProtos.GetAttributeResponseMsg attributesResponse, TransportProtos.SessionInfoProto sessionInfo) {
try {
- LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2MClient(sessionInfo);
- attributesResponse.getSharedAttributeListList().forEach(attr -> {
- String path = this.getPathAttributeUpdate(sessionInfo, attr.getKv().getKey());
- if (path != null) {
- // #1.1
- if (lwM2MClient.getDelayedRequests().containsKey(path) && attr.getTs() > lwM2MClient.getDelayedRequests().get(path).getTs()) {
- lwM2MClient.getDelayedRequests().put(path, attr);
- } else {
- lwM2MClient.getDelayedRequests().put(path, attr);
- }
- }
- });
- // #2.1
- lwM2MClient.getDelayedRequests().forEach((k, v) -> {
- ArrayList listV = new ArrayList<>();
- listV.add(v.getKv());
- this.putDelayedUpdateResourcesClient(lwM2MClient, this.getResourceValueToString(lwM2MClient, k), getJsonObject(listV).get(v.getKv().getKey()), k);
- });
+ List tsKvProtos = attributesResponse.getSharedAttributeListList();
+ this.updateAttriuteFromThingsboard(tsKvProtos, sessionInfo);
} catch (Exception e) {
log.error(String.valueOf(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
+ * => send to client Request Update of value (new value from shared attribute)
+ * and LwM2MClient.delayedRequests.add(path)
+ * #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 updateAttriuteFromThingsboard(List tsKvProtos, TransportProtos.SessionInfoProto sessionInfo) {
+ LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2MClient(sessionInfo);
+ tsKvProtos.forEach(tsKvProto -> {
+ String pathIdVer = this.validatePathIntoProfile(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);
+ }
+ }
+ });
+ // #2.1
+ lwM2MClient.getDelayedRequests().forEach((pathIdVer, tsKvProto) -> {
+ this.updateResourcesValueToClient(lwM2MClient, this.getResourceValueFormatKv(lwM2MClient, pathIdVer),
+ this.lwM2mTransportContextServer.getValueFromKvProto(tsKvProto.getKv()), pathIdVer);
+ });
+ }
+
/**
* @param lwM2MClient -
* @return SessionInfoProto -
@@ -1141,11 +1274,11 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
return new ArrayList<>(namesIsWritable);
}
- private boolean validateResourceInModel(LwM2mClient lwM2mClient, String pathKey, boolean isWritable) {
- ResourceModel resourceModel = lwM2mClient.getResourceModel(pathKey);
- Integer objectId = validateObjectIdFromKey(pathKey);
- String objectVer = validateObjectVerFromKey(pathKey);
- return resourceModel != null && (isWritable ?
+ private boolean validateResourceInModel(LwM2mClient lwM2mClient, String pathIdVer, boolean isWritableNotOptional) {
+ ResourceModel resourceModel = lwM2mClient.getResourceModel(pathIdVer);
+ Integer objectId = new LwM2mPath(convertPathFromIdVerToObjectId(pathIdVer)).getObjectId();
+ String objectVer = validateObjectVerFromKey(pathIdVer);
+ return resourceModel != null && (isWritableNotOptional ?
objectId != null && objectVer != null && objectVer.equals(lwM2mClient.getRegistration().getSupportedVersion(objectId)) && resourceModel.operations.isWritable() :
objectId != null && objectVer != null && objectVer.equals(lwM2mClient.getRegistration().getSupportedVersion(objectId)));
}
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 bdb783e3ff..5e736d8e0f 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
@@ -25,18 +25,21 @@ import org.eclipse.leshan.server.registration.Registration;
import org.eclipse.leshan.server.security.SecurityInfo;
import org.thingsboard.server.gen.transport.TransportProtos;
import org.thingsboard.server.gen.transport.TransportProtos.ValidateDeviceCredentialsResponseMsg;
+import org.thingsboard.server.transport.lwm2m.server.LwM2mQueuedRequest;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportServiceImpl;
import java.util.List;
import java.util.Map;
+import java.util.Queue;
import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.CopyOnWriteArrayList;
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.LwM2mTransportHandler.convertToObjectIdFromIdVer;
+import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromIdVerToObjectId;
@Slf4j
@Data
@@ -54,6 +57,7 @@ public class LwM2mClient implements Cloneable {
private final Map resources;
private final Map delayedRequests;
private final List pendingRequests;
+ private final Queue queuedRequests;
private boolean init;
public Object clone() throws CloneNotSupportedException {
@@ -71,6 +75,7 @@ public class LwM2mClient implements Cloneable {
this.profileId = profileId;
this.sessionId = sessionId;
this.init = false;
+ this.queuedRequests = new ConcurrentLinkedQueue<>();
}
public boolean saveResourceValue(String pathRez, LwM2mResource rez, LwM2mModelProvider modelProvider) {
@@ -78,7 +83,7 @@ public class LwM2mClient implements Cloneable {
this.resources.get(pathRez).setLwM2mResource(rez);
return true;
} else {
- LwM2mPath pathIds = new LwM2mPath(convertToObjectIdFromIdVer(pathRez));
+ LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(pathRez));
ResourceModel resourceModel = modelProvider.getObjectModel(registration).getResourceModel(pathIds.getObjectId(), pathIds.getResourceId());
if (resourceModel != null) {
this.resources.put(pathRez, new ResourceValue(rez, resourceModel));
@@ -103,9 +108,9 @@ public class LwM2mClient implements Cloneable {
* @param modelProvider -
*/
public void deleteResources(String pathIdVer, LwM2mModelProvider modelProvider) {
- Set key = getKeysEqualsIdVer(pathIdVer);
+ Set key = getKeysEqualsIdVer(pathIdVer);
key.forEach(pathRez -> {
- LwM2mPath pathIds = new LwM2mPath(convertToObjectIdFromIdVer(pathRez.toString()));
+ LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(pathRez));
ResourceModel resourceModel = modelProvider.getObjectModel(registration).getResourceModel(pathIds.getObjectId(), pathIds.getResourceId());
if (resourceModel != null) {
this.resources.get(pathRez).setResourceModel(resourceModel);
@@ -122,17 +127,17 @@ public class LwM2mClient implements Cloneable {
* @param modelProvider -
*/
public void updateResourceModel(String idVer, LwM2mModelProvider modelProvider) {
- Set key = getKeysEqualsIdVer(idVer);
- key.forEach(k -> this.saveResourceModel(k.toString(), modelProvider));
+ Set key = getKeysEqualsIdVer(idVer);
+ key.forEach(k -> this.saveResourceModel(k, modelProvider));
}
private void saveResourceModel(String pathRez, LwM2mModelProvider modelProvider) {
- LwM2mPath pathIds = new LwM2mPath(convertToObjectIdFromIdVer(pathRez));
+ LwM2mPath pathIds = new LwM2mPath(convertPathFromIdVerToObjectId(pathRez));
ResourceModel resourceModel = modelProvider.getObjectModel(registration).getResourceModel(pathIds.getObjectId(), pathIds.getResourceId());
this.resources.get(pathRez).setResourceModel(resourceModel);
}
- private Set getKeysEqualsIdVer(String idVer) {
+ private Set getKeysEqualsIdVer(String idVer) {
return this.resources.keySet()
.stream()
.filter(e -> idVer.equals(e.split(LWM2M_SEPARATOR_PATH)[1]))
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java
index 2aea3fdc3e..358f760fef 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContext.java
@@ -20,6 +20,7 @@ import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.gen.transport.TransportProtos;
import java.util.Map;
+import java.util.Set;
import java.util.UUID;
public interface LwM2mClientContext {
@@ -51,4 +52,6 @@ public interface LwM2mClientContext {
Map setProfiles(Map profiles);
boolean addUpdateProfileParameters(DeviceProfile deviceProfile);
+
+ Set getSupportedIdVerInClient(Registration registration);
}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
index 926d12516b..cca729f5ff 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientContextImpl.java
@@ -15,6 +15,7 @@
*/
package org.thingsboard.server.transport.lwm2m.server.client;
+import org.eclipse.leshan.core.node.LwM2mPath;
import org.eclipse.leshan.server.registration.Registration;
import org.eclipse.leshan.server.security.EditableSecurityStore;
import org.springframework.stereotype.Service;
@@ -27,11 +28,14 @@ import org.thingsboard.server.transport.lwm2m.secure.ReadResultSecurityStore;
import org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler;
import org.thingsboard.server.transport.lwm2m.utils.TypeServer;
+import java.util.Arrays;
import java.util.Map;
+import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import static org.thingsboard.server.transport.lwm2m.secure.LwM2MSecurityMode.NO_SEC;
+import static org.thingsboard.server.transport.lwm2m.server.LwM2mTransportHandler.convertPathFromObjectIdToIdVer;
@Service
@TbLwM2mTransportComponent
@@ -90,7 +94,7 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
@Override
public LwM2mClient updateInSessionsLwM2MClient(Registration registration) {
if (this.lwM2mClients.get(registration.getEndpoint()) == null) {
- addLwM2mClientToSession(registration.getEndpoint());
+ this.addLwM2mClientToSession(registration.getEndpoint());
}
LwM2mClient lwM2MClient = lwM2mClients.get(registration.getEndpoint());
lwM2MClient.setRegistration(registration);
@@ -169,4 +173,21 @@ public class LwM2mClientContextImpl implements LwM2mClientContext {
}
return false;
}
+
+ /**
+ * if isVer - ok or default ver=DEFAULT_LWM2M_VERSION
+ * @param registration -
+ * @return - all objectIdVer in client
+ */
+ @Override
+ public Set getSupportedIdVerInClient(Registration registration) {
+ Set clientObjects = ConcurrentHashMap.newKeySet();
+ Arrays.stream(registration.getObjectLinks()).forEach(url -> {
+ LwM2mPath pathIds = new LwM2mPath(url.getUrl());
+ if (!pathIds.isRoot()) {
+ clientObjects.add(convertPathFromObjectIdToIdVer(url.getUrl(), registration));
+ }
+ });
+ return (clientObjects.size() > 0) ? clientObjects : null;
+ }
}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientProfile.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientProfile.java
index 1c4042bd1a..19af453f3c 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientProfile.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/LwM2mClientProfile.java
@@ -56,6 +56,13 @@ public class LwM2mClientProfile {
*/
private JsonArray postObserveProfile;
+ /**
+ * "attributeLwm2m": {"/3_1.0": {"ver": "currentTimeTest11"},
+ * "/3_1.0/0": {"gt": 17},
+ * "/3_1.0/0/9": {"pmax": 45}, "/3_1.2": {ver": "3_1.2"}}
+ */
+ private JsonObject postAttributeLwm2mProfile;
+
public LwM2mClientProfile clone() {
LwM2mClientProfile lwM2mClientProfile = new LwM2mClientProfile();
lwM2mClientProfile.postClientLwM2mSettings = this.deepCopy(this.postClientLwM2mSettings, JsonObject.class);
@@ -63,6 +70,7 @@ public class LwM2mClientProfile {
lwM2mClientProfile.postAttributeProfile = this.deepCopy(this.postAttributeProfile, JsonArray.class);
lwM2mClientProfile.postTelemetryProfile = this.deepCopy(this.postTelemetryProfile, JsonArray.class);
lwM2mClientProfile.postObserveProfile = this.deepCopy(this.postObserveProfile, JsonArray.class);
+ lwM2mClientProfile.postAttributeLwm2mProfile = this.deepCopy(this.postAttributeLwm2mProfile, JsonObject.class);
return lwM2mClientProfile;
}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsAddKeyValueProto.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsAddKeyValueProto.java
new file mode 100644
index 0000000000..5c7edbe328
--- /dev/null
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsAddKeyValueProto.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.client;
+
+import lombok.Data;
+import org.thingsboard.server.gen.transport.TransportProtos;
+
+import java.util.ArrayList;
+import java.util.List;
+
+@Data
+public class ResultsAddKeyValueProto {
+ List resultAttributes;
+ List resultTelemetries;
+
+ public ResultsAddKeyValueProto() {
+ this.resultAttributes = new ArrayList<>();
+ this.resultTelemetries = new ArrayList<>();
+ }
+
+}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.java
new file mode 100644
index 0000000000..36cd8943fb
--- /dev/null
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/server/client/ResultsResourceValue.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.client;
+
+import lombok.Data;
+import org.thingsboard.server.common.data.kv.DataType;
+
+@Data
+public class ResultsResourceValue {
+ DataType dataType;
+ Object value;
+ String resourceName;
+
+ public ResultsResourceValue (DataType dataType, Object value, String resourceName) {
+ this.dataType = dataType;
+ this.value = value;
+ this.resourceName = resourceName;
+ }
+}
diff --git a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2mValueConverterImpl.java b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2mValueConverterImpl.java
index 0faa20952e..dce34c9cf6 100644
--- a/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2mValueConverterImpl.java
+++ b/common/transport/lwm2m/src/main/java/org/thingsboard/server/transport/lwm2m/utils/LwM2mValueConverterImpl.java
@@ -18,14 +18,12 @@ package org.thingsboard.server.transport.lwm2m.utils;
import lombok.extern.slf4j.Slf4j;
import org.eclipse.leshan.core.model.ResourceModel.Type;
import org.eclipse.leshan.core.node.LwM2mPath;
+import org.eclipse.leshan.core.node.ObjectLink;
import org.eclipse.leshan.core.node.codec.CodecException;
import org.eclipse.leshan.core.node.codec.LwM2mValueConverter;
import org.eclipse.leshan.core.util.Hex;
import org.eclipse.leshan.core.util.StringUtils;
-import javax.xml.datatype.DatatypeConfigurationException;
-import javax.xml.datatype.DatatypeFactory;
-import javax.xml.datatype.XMLGregorianCalendar;
import java.math.BigInteger;
import java.text.DateFormat;
import java.text.SimpleDateFormat;
@@ -111,15 +109,16 @@ public class LwM2mValueConverterImpl implements LwM2mValueConverter {
case INTEGER:
log.debug("Trying to convert long value {} to date", value);
/** let's assume we received the millisecond since 1970/1/1 */
- return new Date((Long) value);
+ return new Date(((Number) value).longValue() * 1000L);
case STRING:
log.debug("Trying to convert string value {} to date", value);
/** let's assume we received an ISO 8601 format date */
try {
- DatatypeFactory datatypeFactory = DatatypeFactory.newInstance();
- XMLGregorianCalendar cal = datatypeFactory.newXMLGregorianCalendar((String) value);
- return cal.toGregorianCalendar().getTime();
- } catch (DatatypeConfigurationException | IllegalArgumentException e) {
+ return new Date(Long.decode(value.toString()));
+// DatatypeFactory datatypeFactory = DatatypeFactory.newInstance();
+// XMLGregorianCalendar cal = datatypeFactory.newXMLGregorianCalendar((String) value);
+// return cal.toGregorianCalendar().getTime();
+ } catch (IllegalArgumentException e) {
log.debug("Unable to convert string to date", e);
throw new CodecException("Unable to convert string (%s) to date for resource %s", value,
resourcePath);
@@ -147,6 +146,8 @@ public class LwM2mValueConverterImpl implements LwM2mValueConverter {
return formatter.format(new Date(timeValue));
case OPAQUE:
return Hex.encodeHexString((byte[])value);
+ case OBJLNK:
+ return ObjectLink.decodeFromString((String) value);
default:
break;
}
@@ -164,10 +165,14 @@ public class LwM2mValueConverterImpl implements LwM2mValueConverter {
}
}
break;
+ case OBJLNK:
+ if (currentType == Type.STRING) {
+ return ObjectLink.fromPath(value.toString());
+ }
default:
}
throw new CodecException("Invalid value type for resource %s, expected %s, got %s", resourcePath, expectedType,
currentType);
}
-}
+ }
diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java
index c948e755c2..30178bd44f 100644
--- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java
+++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/adaptors/ProtoMqttAdaptor.java
@@ -15,6 +15,7 @@
*/
package org.thingsboard.server.transport.mqtt.adaptors;
+import com.google.gson.JsonElement;
import com.google.gson.JsonParser;
import com.google.protobuf.Descriptors;
import com.google.protobuf.DynamicMessage;
@@ -64,9 +65,9 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
public TransportProtos.PostAttributeMsg convertToPostAttributes(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx;
byte[] bytes = toBytes(inbound.payload());
- Descriptors.Descriptor attributesDynamicMessage = getDescriptor(deviceSessionCtx.getAttributesDynamicMessageDescriptor());
+ Descriptors.Descriptor attributesDynamicMessageDescriptor = getDescriptor(deviceSessionCtx.getAttributesDynamicMessageDescriptor());
try {
- return JsonConverter.convertToAttributesProto(new JsonParser().parse(dynamicMsgToJson(bytes, attributesDynamicMessage)));
+ return JsonConverter.convertToAttributesProto(new JsonParser().parse(dynamicMsgToJson(bytes, attributesDynamicMessageDescriptor)));
} catch (Exception e) {
throw new AdaptorException(e);
}
@@ -86,8 +87,8 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
public TransportProtos.GetAttributeRequestMsg convertToGetAttributes(MqttDeviceAwareSessionContext ctx, MqttPublishMessage inbound) throws AdaptorException {
byte[] bytes = toBytes(inbound.payload());
String topicName = inbound.variableHeader().topicName();
- int requestId = getRequestId(topicName, MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX);
try {
+ int requestId = getRequestId(topicName, MqttTopics.DEVICE_ATTRIBUTES_REQUEST_TOPIC_PREFIX);
return ProtoConverter.convertToGetAttributeRequestMessage(bytes, requestId);
} catch (InvalidProtocolBufferException e) {
log.warn("Failed to decode get attributes request", e);
@@ -97,10 +98,15 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
@Override
public TransportProtos.ToDeviceRpcResponseMsg convertToDeviceRpcResponse(MqttDeviceAwareSessionContext ctx, MqttPublishMessage mqttMsg) throws AdaptorException {
+ DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx;
+ String topicName = mqttMsg.variableHeader().topicName();
byte[] bytes = toBytes(mqttMsg.payload());
+ Descriptors.Descriptor rpcResponseDynamicMessageDescriptor = getDescriptor(deviceSessionCtx.getRpcResponseDynamicMessageDescriptor());
try {
- return TransportProtos.ToDeviceRpcResponseMsg.parseFrom(bytes);
- } catch (RuntimeException | InvalidProtocolBufferException e) {
+ int requestId = getRequestId(topicName, MqttTopics.DEVICE_RPC_RESPONSE_TOPIC);
+ JsonElement response = new JsonParser().parse(dynamicMsgToJson(bytes, rpcResponseDynamicMessageDescriptor));
+ return TransportProtos.ToDeviceRpcResponseMsg.newBuilder().setRequestId(requestId).setPayload(response.toString()).build();
+ } catch (Exception e) {
log.warn("Failed to decode Rpc response", e);
throw new AdaptorException(e);
}
@@ -145,8 +151,10 @@ public class ProtoMqttAdaptor implements MqttTransportAdaptor {
@Override
- public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) {
- return Optional.of(createMqttPublishMsg(ctx, MqttTopics.DEVICE_RPC_REQUESTS_TOPIC + rpcRequest.getRequestId(), ProtoConverter.convertToRpcRequest(rpcRequest)));
+ public Optional convertToPublish(MqttDeviceAwareSessionContext ctx, TransportProtos.ToDeviceRpcRequestMsg rpcRequest) throws AdaptorException {
+ DeviceSessionCtx deviceSessionCtx = (DeviceSessionCtx) ctx;
+ DynamicMessage.Builder rpcRequestDynamicMessageBuilder = deviceSessionCtx.getRpcRequestDynamicMessageBuilder();
+ return Optional.of(createMqttPublishMsg(ctx, MqttTopics.DEVICE_RPC_REQUESTS_TOPIC + rpcRequest.getRequestId(), ProtoConverter.convertToRpcRequest(rpcRequest, rpcRequestDynamicMessageBuilder)));
}
@Override
diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
index efc5c057ed..3804d96cf6 100644
--- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
+++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/session/DeviceSessionCtx.java
@@ -16,6 +16,7 @@
package org.thingsboard.server.transport.mqtt.session;
import com.google.protobuf.Descriptors;
+import com.google.protobuf.DynamicMessage;
import io.netty.channel.ChannelHandlerContext;
import lombok.Getter;
import lombok.Setter;
@@ -60,6 +61,8 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
private volatile TransportPayloadType payloadType = TransportPayloadType.JSON;
private volatile Descriptors.Descriptor attributesDynamicMessageDescriptor;
private volatile Descriptors.Descriptor telemetryDynamicMessageDescriptor;
+ private volatile Descriptors.Descriptor rpcResponseDynamicMessageDescriptor;
+ private volatile DynamicMessage.Builder rpcRequestDynamicMessageBuilder;
@Getter
@Setter
@@ -102,6 +105,14 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
return attributesDynamicMessageDescriptor;
}
+ public Descriptors.Descriptor getRpcResponseDynamicMessageDescriptor() {
+ return rpcResponseDynamicMessageDescriptor;
+ }
+
+ public DynamicMessage.Builder getRpcRequestDynamicMessageBuilder() {
+ return rpcRequestDynamicMessageBuilder;
+ }
+
@Override
public void setDeviceProfile(DeviceProfile deviceProfile) {
super.setDeviceProfile(deviceProfile);
@@ -136,5 +147,7 @@ public class DeviceSessionCtx extends MqttDeviceAwareSessionContext {
ProtoTransportPayloadConfiguration protoTransportPayloadConfig = (ProtoTransportPayloadConfiguration) transportPayloadTypeConfiguration;
telemetryDynamicMessageDescriptor = protoTransportPayloadConfig.getTelemetryDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceTelemetryProtoSchema());
attributesDynamicMessageDescriptor = protoTransportPayloadConfig.getAttributesDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceAttributesProtoSchema());
+ rpcResponseDynamicMessageDescriptor = protoTransportPayloadConfig.getRpcResponseDynamicMessageDescriptor(protoTransportPayloadConfig.getDeviceRpcResponseProtoSchema());
+ rpcRequestDynamicMessageBuilder = protoTransportPayloadConfig.getRpcRequestDynamicMessageBuilder(protoTransportPayloadConfig.getDeviceRpcRequestProtoSchema());
}
}
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java
index 42c99bc428..d0d471a784 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/AdaptorException.java
@@ -31,4 +31,8 @@ public class AdaptorException extends Exception {
super(cause);
}
+ public AdaptorException(String message, Exception cause) {
+ super(message, cause);
+ }
+
}
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java
index ed440ebcfe..620728cda4 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/adaptor/ProtoConverter.java
@@ -15,10 +15,13 @@
*/
package org.thingsboard.server.common.transport.adaptor;
+import com.google.gson.Gson;
import com.google.gson.JsonElement;
+import com.google.gson.JsonObject;
import com.google.gson.JsonParser;
-import com.google.gson.JsonPrimitive;
+import com.google.protobuf.DynamicMessage;
import com.google.protobuf.InvalidProtocolBufferException;
+import com.google.protobuf.util.JsonFormat;
import lombok.extern.slf4j.Slf4j;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
@@ -34,6 +37,7 @@ import java.util.List;
@Slf4j
public class ProtoConverter {
+ public static final Gson GSON = new Gson();
public static final JsonParser JSON_PARSER = new JsonParser();
public static TransportProtos.PostTelemetryMsg convertToTelemetryProto(byte[] payload) throws InvalidProtocolBufferException, IllegalArgumentException {
@@ -170,26 +174,20 @@ public class ProtoConverter {
return kvList;
}
- public static byte[] convertToRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRpcRequestMsg) {
- TransportProtos.ToDeviceRpcRequestMsg.Builder toDeviceRpcRequestMsgBuilder = toDeviceRpcRequestMsg.newBuilderForType();
- toDeviceRpcRequestMsgBuilder.mergeFrom(toDeviceRpcRequestMsg);
- toDeviceRpcRequestMsgBuilder.setParams(parseParams(toDeviceRpcRequestMsg));
- TransportProtos.ToDeviceRpcRequestMsg result = toDeviceRpcRequestMsgBuilder.build();
- return result.toByteArray();
- }
-
- private static String parseParams(TransportProtos.ToDeviceRpcRequestMsg toDeviceRpcRequestMsg) {
+ public static byte[] convertToRpcRequest(TransportProtos.ToDeviceRpcRequestMsg toDeviceRpcRequestMsg, DynamicMessage.Builder rpcRequestDynamicMessageBuilder) throws AdaptorException {
+ rpcRequestDynamicMessageBuilder.clear();
+ JsonObject rpcRequestJson = new JsonObject();
+ rpcRequestJson.addProperty("method", toDeviceRpcRequestMsg.getMethodName());
+ rpcRequestJson.addProperty("requestId", toDeviceRpcRequestMsg.getRequestId());
String params = toDeviceRpcRequestMsg.getParams();
- JsonElement jsonElementParams = JSON_PARSER.parse(params);
- if (!jsonElementParams.isJsonPrimitive()) {
- return params;
- } else {
- JsonPrimitive primitiveParams = jsonElementParams.getAsJsonPrimitive();
- if (jsonElementParams.getAsJsonPrimitive().isString()) {
- return primitiveParams.getAsString();
- } else {
- return params;
- }
+ try {
+ JsonElement paramsElement = JSON_PARSER.parse(params);
+ rpcRequestJson.add("params", paramsElement);
+ JsonFormat.parser().ignoringUnknownFields().merge(GSON.toJson(rpcRequestJson), rpcRequestDynamicMessageBuilder);
+ DynamicMessage dynamicRpcRequest = rpcRequestDynamicMessageBuilder.build();
+ return dynamicRpcRequest.toByteArray();
+ } catch (Exception e) {
+ throw new AdaptorException("Failed to convert ToDeviceRpcRequestMsg to Dynamic Rpc request message due to: ", e);
}
}
}
diff --git a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/lwm2m/LwM2MTransportConfigServer.java b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/lwm2m/LwM2MTransportConfigServer.java
index 5603721474..cb2af9eb3f 100644
--- a/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/lwm2m/LwM2MTransportConfigServer.java
+++ b/common/transport/transport-api/src/main/java/org/thingsboard/server/common/transport/lwm2m/LwM2MTransportConfigServer.java
@@ -207,6 +207,8 @@ public class LwM2MTransportConfigServer {
FULL_FILE_PATH = Paths.get(BASE_DIR_PATH.replaceAll("bin$", ""));
} else if (BASE_DIR_PATH.endsWith("conf")) {
FULL_FILE_PATH = Paths.get(BASE_DIR_PATH.replaceAll("conf$", ""));
+ } else if (BASE_DIR_PATH.endsWith("application")) {
+ FULL_FILE_PATH = Paths.get(BASE_DIR_PATH.substring(0, BASE_DIR_PATH.length() - "application".length()));
} else {
FULL_FILE_PATH = Paths.get(BASE_DIR_PATH);
}
diff --git a/dao/pom.xml b/dao/pom.xml
index 9ba3033970..cc25e1a13c 100644
--- a/dao/pom.xml
+++ b/dao/pom.xml
@@ -227,10 +227,6 @@
org.elasticsearch.client
rest