resourceUpdateMsgOpt) {
String idVer = resourceUpdateMsgOpt.get().getResourceKey();
lwM2mClientContext.getLwM2mClients().values().stream().forEach(e -> e.updateResourceModel(idVer, this.lwM2mTransportContextServer.getLwM2MTransportConfigServer().getModelProvider()));
}
/**
- *
* @param resourceDeleteMsgOpt -
*/
@Override
@@ -374,6 +385,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"
*
@@ -417,6 +436,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.
}
@@ -485,15 +505,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);
}
/**
@@ -600,41 +625,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
@@ -654,21 +699,6 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
lwM2MClient.setProfileId(device.getDeviceProfileId().getId());
}
- /**
- * @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
@@ -713,16 +743,18 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
private void addParameters(String path, JsonObject parameters, Registration registration) {
LwM2mClient lwM2MClient = lwM2mClientContext.getLwM2mClientWithReg(registration, null);
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) {
+ if (names != null && names.has(path)) {
+ String resName = names.get(path).getAsString();
+ if (resName != null && !resName.isEmpty()) {
+ try {
+ String resValue = this.getResourceValueToString(lwM2MClient, path);
parameters.addProperty(resName, resValue);
+ } 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: [{}]", path, names);
}
}
@@ -739,7 +771,7 @@ public class LwM2mTransportServiceImpl implements LwM2mTransportService {
/**
* @param lwM2MClient -
- * @param path -
+ * @param path -
* @return - return value of Resource by idPath
*/
private LwM2mResource returnResourceValueFromLwM2MClient(LwM2mClient lwM2MClient, String path) {
@@ -768,6 +800,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
@@ -778,6 +811,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 -
@@ -792,6 +828,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();
@@ -800,32 +837,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
@@ -859,10 +905,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());
+ }
});
}
}
@@ -922,7 +972,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()
@@ -932,6 +982,71 @@ 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 -> {
@@ -1140,11 +1255,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(convertToObjectIdFromIdVer(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 186f085397..56e8fd5c7b 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,13 +25,16 @@ 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;
@@ -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) {
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..1ce16b9b92 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.convertToIdVerFromObjectId;
@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(convertToIdVerFromObjectId(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/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
-
- org.eclipse.leshan
- leshan-core
-
diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
index 4cc853b2b9..91bfa954fd 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/event/BaseEventService.java
@@ -23,6 +23,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.Event;
+import org.thingsboard.server.common.data.event.EventFilter;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
@@ -30,7 +31,6 @@ import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.exception.DataValidationException;
import org.thingsboard.server.dao.service.DataValidator;
-import java.nio.charset.StandardCharsets;
import java.util.List;
import java.util.Optional;
@@ -111,6 +111,11 @@ public class BaseEventService implements EventService {
return eventDao.findLatestEvents(tenantId.getId(), entityId, eventType, limit);
}
+ @Override
+ public PageData findEventsByFilter(TenantId tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink) {
+ return eventDao.findEventByFilter(tenantId.getId(), entityId, eventFilter, pageLink);
+ }
+
@Override
public void removeEvents(TenantId tenantId, EntityId entityId) {
PageData eventPageData;
diff --git a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
index 0cb4fccfb3..ba4e86c95e 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
@@ -17,6 +17,7 @@ package org.thingsboard.server.dao.event;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.Event;
+import org.thingsboard.server.common.data.event.EventFilter;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageData;
@@ -88,6 +89,8 @@ public interface EventDao extends Dao {
*/
PageData findEvents(UUID tenantId, EntityId entityId, String eventType, TimePageLink pageLink);
+ PageData findEventByFilter(UUID tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink);
+
/**
* Find latest events by tenantId, entityId and eventType.
*
diff --git a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseTbResourceService.java b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java
similarity index 55%
rename from dao/src/main/java/org/thingsboard/server/dao/resource/BaseTbResourceService.java
rename to dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java
index 0df2f19e80..97675d4b3e 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/resource/BaseTbResourceService.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/resource/BaseResourceService.java
@@ -17,10 +17,6 @@ package org.thingsboard.server.dao.resource;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
-import org.eclipse.leshan.core.model.DDFFileParser;
-import org.eclipse.leshan.core.model.DefaultDDFFileValidator;
-import org.eclipse.leshan.core.model.InvalidDDFFileException;
-import org.eclipse.leshan.core.model.ObjectModel;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.stereotype.Service;
import org.thingsboard.server.common.data.ResourceType;
@@ -29,9 +25,6 @@ import org.thingsboard.server.common.data.TbResourceInfo;
import org.thingsboard.server.common.data.Tenant;
import org.thingsboard.server.common.data.id.TbResourceId;
import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.common.data.lwm2m.LwM2mInstance;
-import org.thingsboard.server.common.data.lwm2m.LwM2mObject;
-import org.thingsboard.server.common.data.lwm2m.LwM2mResourceObserve;
import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.PageLink;
import org.thingsboard.server.dao.exception.DataValidationException;
@@ -41,63 +34,30 @@ import org.thingsboard.server.dao.service.PaginatedRemover;
import org.thingsboard.server.dao.service.Validator;
import org.thingsboard.server.dao.tenant.TenantDao;
-import java.io.ByteArrayInputStream;
-import java.io.IOException;
-import java.util.ArrayList;
-import java.util.Base64;
-import java.util.Comparator;
import java.util.List;
import java.util.Optional;
-import java.util.stream.Collectors;
-import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_KEY;
-import static org.thingsboard.server.common.data.lwm2m.LwM2mConstants.LWM2M_SEPARATOR_SEARCH_TEXT;
import static org.thingsboard.server.dao.device.DeviceServiceImpl.INCORRECT_TENANT_ID;
import static org.thingsboard.server.dao.service.Validator.validateId;
@Service
@Slf4j
-public class BaseTbResourceService implements TbResourceService {
+public class BaseResourceService implements ResourceService {
public static final String INCORRECT_RESOURCE_ID = "Incorrect resourceId ";
private final TbResourceDao resourceDao;
private final TbResourceInfoDao resourceInfoDao;
private final TenantDao tenantDao;
- private final DDFFileParser ddfFileParser;
- public BaseTbResourceService(TbResourceDao resourceDao, TbResourceInfoDao resourceInfoDao, TenantDao tenantDao) {
+
+ public BaseResourceService(TbResourceDao resourceDao, TbResourceInfoDao resourceInfoDao, TenantDao tenantDao) {
this.resourceDao = resourceDao;
this.resourceInfoDao = resourceInfoDao;
this.tenantDao = tenantDao;
- this.ddfFileParser = new DDFFileParser(new DefaultDDFFileValidator());
}
@Override
- public TbResource saveResource(TbResource resource) throws InvalidDDFFileException, IOException {
- log.trace("Executing saveResource [{}]", resource);
- if (StringUtils.isEmpty(resource.getData())) {
- throw new DataValidationException("Resource data should be specified!");
- }
- if (ResourceType.LWM2M_MODEL.equals(resource.getResourceType())) {
- List objectModels =
- ddfFileParser.parseEx(new ByteArrayInputStream(Base64.getDecoder().decode(resource.getData())), resource.getSearchText());
- if (!objectModels.isEmpty()) {
- ObjectModel objectModel = objectModels.get(0);
-
- String resourceKey = objectModel.id + LWM2M_SEPARATOR_KEY + objectModel.getVersion();
- String name = objectModel.name;
- resource.setResourceKey(resourceKey);
- if (resource.getId() == null) {
- resource.setTitle(name + " id=" + objectModel.id + " v" + objectModel.getVersion());
- }
- resource.setSearchText(resourceKey + LWM2M_SEPARATOR_SEARCH_TEXT + name);
- } else {
- throw new DataValidationException(String.format("Could not parse the XML of objectModel with name %s", resource.getSearchText()));
- }
- } else {
- resource.setResourceKey(resource.getFileName());
- }
-
+ public TbResource saveResource(TbResource resource) {
resourceValidator.validate(resource, TbResourceInfo::getTenantId);
try {
@@ -111,7 +71,6 @@ public class BaseTbResourceService implements TbResourceService {
throw t;
}
}
-
}
@Override
@@ -156,31 +115,17 @@ public class BaseTbResourceService implements TbResourceService {
}
@Override
- public List findLwM2mObjectPage(TenantId tenantId, String sortProperty, String sortOrder, PageLink pageLink) {
- log.trace("Executing findByTenantId [{}]", tenantId);
+ public List findTenantResourcesByResourceTypeAndObjectIds(TenantId tenantId, ResourceType resourceType, String[] objectIds) {
+ log.trace("Executing findTenantResourcesByResourceTypeAndObjectIds [{}][{}][{}]", tenantId, resourceType, objectIds);
validateId(tenantId, INCORRECT_TENANT_ID + tenantId);
- PageData resourcePageData = resourceDao.findResourcesByTenantIdAndResourceType(
- tenantId,
- ResourceType.LWM2M_MODEL, pageLink);
- return resourcePageData.getData().stream()
- .map(this::toLwM2mObject)
- .sorted(getComparator(sortProperty, sortOrder))
- .collect(Collectors.toList());
+ return resourceDao.findResourcesByTenantIdAndResourceType(tenantId, resourceType, objectIds, null);
}
@Override
- public List findLwM2mObject(TenantId tenantId, String sortOrder,
- String sortProperty,
- String[] objectIds) {
- log.trace("Executing findByTenantId [{}]", tenantId);
+ public PageData findTenantResourcesByResourceTypeAndPageLink(TenantId tenantId, ResourceType resourceType, PageLink pageLink) {
+ log.trace("Executing findTenantResourcesByResourceTypeAndPageLink [{}][{}][{}]", tenantId, resourceType, pageLink);
validateId(tenantId, INCORRECT_TENANT_ID + tenantId);
- List resources = resourceDao.findResourcesByTenantIdAndResourceType(tenantId, ResourceType.LWM2M_MODEL,
- objectIds,
- null);
- return resources.stream()
- .map(this::toLwM2mObject)
- .sorted(getComparator(sortProperty, sortOrder))
- .collect(Collectors.toList());
+ return resourceDao.findResourcesByTenantIdAndResourceType(tenantId, resourceType, pageLink);
}
@Override
@@ -190,50 +135,6 @@ public class BaseTbResourceService implements TbResourceService {
tenantResourcesRemover.removeEntities(tenantId, tenantId);
}
- private LwM2mObject toLwM2mObject(TbResource resource) {
- try {
- DDFFileParser ddfFileParser = new DDFFileParser(new DefaultDDFFileValidator());
- List objectModels =
- ddfFileParser.parseEx(new ByteArrayInputStream(Base64.getDecoder().decode(resource.getData())), resource.getSearchText());
- if (objectModels.size() == 0) {
- return null;
- } else {
- ObjectModel obj = objectModels.get(0);
- LwM2mObject lwM2mObject = new LwM2mObject();
- lwM2mObject.setId(obj.id);
- lwM2mObject.setKeyId(resource.getResourceKey());
- lwM2mObject.setName(obj.name);
- lwM2mObject.setMultiple(obj.multiple);
- lwM2mObject.setMandatory(obj.mandatory);
- LwM2mInstance instance = new LwM2mInstance();
- instance.setId(0);
- List resources = new ArrayList<>();
- obj.resources.forEach((k, v) -> {
- if (!v.operations.isExecutable()) {
- LwM2mResourceObserve lwM2MResourceObserve = new LwM2mResourceObserve(k, v.name, false, false, false);
- resources.add(lwM2MResourceObserve);
- }
- });
- instance.setResources(resources.toArray(LwM2mResourceObserve[]::new));
- lwM2mObject.setInstances(new LwM2mInstance[]{instance});
- return lwM2mObject;
- }
- } catch (IOException | InvalidDDFFileException e) {
- log.error("Could not parse the XML of objectModel with name [{}]", resource.getSearchText(), e);
- return null;
- }
- }
-
- private Comparator super LwM2mObject> getComparator(String sortProperty, String sortOrder) {
- Comparator comparator;
- if ("name".equals(sortProperty)) {
- comparator = Comparator.comparing(LwM2mObject::getName);
- } else {
- comparator = Comparator.comparingLong(LwM2mObject::getId);
- }
- return "DESC".equals(sortOrder) ? comparator.reversed() : comparator;
- }
-
private DataValidator resourceValidator = new DataValidator<>() {
@Override
@@ -259,9 +160,6 @@ public class BaseTbResourceService implements TbResourceService {
throw new DataValidationException("Resource is referencing to non-existent tenant!");
}
}
- if (resource.getResourceType().equals(ResourceType.LWM2M_MODEL) && toLwM2mObject(resource) == null) {
- throw new DataValidationException(String.format("Could not parse the XML of objectModel with name %s", resource.getSearchText()));
- }
}
};
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/component/AbstractComponentDescriptorInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/component/AbstractComponentDescriptorInsertRepository.java
index c6c545b574..f4eb71eaed 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sql/component/AbstractComponentDescriptorInsertRepository.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sql/component/AbstractComponentDescriptorInsertRepository.java
@@ -51,11 +51,11 @@ public abstract class AbstractComponentDescriptorInsertRepository implements Com
TransactionStatus transaction = getTransactionStatus(TransactionDefinition.PROPAGATION_REQUIRES_NEW);
try {
componentDescriptorEntity = processSaveOrUpdate(entity, insertOrUpdateOnUniqueKeyConflict);
+ transactionManager.commit(transaction);
} catch (Throwable th) {
log.trace("Could not execute the update statement for Component Descriptor with id {}, name {} and entityType {}", entity.getUuid(), entity.getName(), entity.getType());
transactionManager.rollback(transaction);
}
- transactionManager.commit(transaction);
} else {
log.trace("Could not execute the insert statement for Component Descriptor with id {}, name {} and entityType {}", entity.getUuid(), entity.getName(), entity.getType());
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java
index fde9498702..93a33d1f9f 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/AbstractEventInsertRepository.java
@@ -51,11 +51,11 @@ public abstract class AbstractEventInsertRepository implements EventInsertReposi
TransactionStatus transaction = getTransactionStatus(TransactionDefinition.PROPAGATION_REQUIRES_NEW);
try {
eventEntity = processSaveOrUpdate(entity, insertOrUpdateOnUniqueKeyConflict);
+ transactionManager.commit(transaction);
} catch (Throwable th) {
log.trace("Could not execute the update statement for Entity with entityId {} and entityType {}", entity.getEventUid(), entity.getEventType());
transactionManager.rollback(transaction);
}
- transactionManager.commit(transaction);
} else {
log.trace("Could not execute the insert statement for Entity with entityId {} and entityType {}", entity.getEventUid(), entity.getEventType());
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventRepository.java
index c68d6f1a7e..270ab26329 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventRepository.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/EventRepository.java
@@ -44,11 +44,11 @@ public interface EventRepository extends PagingAndSortingRepository findLatestByTenantIdAndEntityTypeAndEntityIdAndEventType(
- @Param("tenantId") UUID tenantId,
- @Param("entityType") EntityType entityType,
- @Param("entityId") UUID entityId,
- @Param("eventType") String eventType,
- Pageable pageable);
+ @Param("tenantId") UUID tenantId,
+ @Param("entityType") EntityType entityType,
+ @Param("entityId") UUID entityId,
+ @Param("eventType") String eventType,
+ Pageable pageable);
@Query("SELECT e FROM EventEntity e WHERE " +
"e.tenantId = :tenantId " +
@@ -80,4 +80,165 @@ public interface EventRepository extends PagingAndSortingRepository= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:type IS NULL OR lower(json_body->>'type') LIKE concat('%', lower(:type\\:\\:varchar), '%')) " +
+ "AND (:server IS NULL OR lower(json_body->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:entityName IS NULL OR lower(json_body->>'entityName') LIKE concat('%', lower(:entityName\\:\\:varchar), '%')) " +
+ "AND (:relationType IS NULL OR lower(json_body->>'relationType') LIKE concat('%', lower(:relationType\\:\\:varchar), '%')) " +
+ "AND (:bodyEntityId IS NULL OR lower(json_body->>'entityId') LIKE concat('%', lower(:bodyEntityId\\:\\:varchar), '%')) " +
+ "AND (:msgType IS NULL OR lower(json_body->>'msgType') LIKE concat('%', lower(:msgType\\:\\:varchar), '%')) " +
+ "AND ((:isError = FALSE) OR (json_body->>'error') IS NOT NULL) " +
+ "AND (:error IS NULL OR lower(json_body->>'error') LIKE concat('%', lower(:error\\:\\:varchar), '%')) " +
+ "AND (:data IS NULL OR lower(json_body->>'data') LIKE concat('%', lower(:data\\:\\:varchar), '%')) " +
+ "AND (:metadata IS NULL OR lower(json_body->>'metadata') LIKE concat('%', lower(:metadata\\:\\:varchar), '%')) ",
+ countQuery = "SELECT count(*) FROM " +
+ "(SELECT *, e.body\\:\\:jsonb as json_body FROM event e WHERE " +
+ "e.tenant_id = :tenantId " +
+ "AND e.entity_type = :entityType " +
+ "AND e.entity_id = :entityId " +
+ "AND e.event_type = :eventType " +
+ "AND e.created_time >= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:type IS NULL OR lower(json_body->>'type') LIKE concat('%', lower(:type\\:\\:varchar), '%')) " +
+ "AND (:server IS NULL OR lower(json_body->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:entityName IS NULL OR lower(json_body->>'entityName') LIKE concat('%', lower(:entityName\\:\\:varchar), '%')) " +
+ "AND (:relationType IS NULL OR lower(json_body->>'relationType') LIKE concat('%', lower(:relationType\\:\\:varchar), '%')) " +
+ "AND (:bodyEntityId IS NULL OR lower(json_body->>'entityId') LIKE concat('%', lower(:bodyEntityId\\:\\:varchar), '%')) " +
+ "AND (:msgType IS NULL OR lower(json_body->>'msgType') LIKE concat('%', lower(:msgType\\:\\:varchar), '%')) " +
+ "AND ((:isError = FALSE) OR (json_body->>'error') IS NOT NULL) " +
+ "AND (:error IS NULL OR lower(json_body->>'error') LIKE concat('%', lower(:error\\:\\:varchar), '%')) " +
+ "AND (:data IS NULL OR lower(json_body->>'data') LIKE concat('%', lower(:data\\:\\:varchar), '%')) " +
+ "AND (:metadata IS NULL OR lower(json_body->>'metadata') LIKE concat('%', lower(:metadata\\:\\:varchar), '%'))"
+ )
+ Page findDebugRuleNodeEvents(@Param("tenantId") UUID tenantId,
+ @Param("entityId") UUID entityId,
+ @Param("entityType") String entityType,
+ @Param("eventType") String eventType,
+ @Param("startTime") Long startTime,
+ @Param("endTime") Long endTime,
+ @Param("type") String type,
+ @Param("server") String server,
+ @Param("entityName") String entityName,
+ @Param("relationType") String relationType,
+ @Param("bodyEntityId") String bodyEntityId,
+ @Param("msgType") String msgType,
+ @Param("isError") boolean isError,
+ @Param("error") String error,
+ @Param("data") String data,
+ @Param("metadata") String metadata,
+ Pageable pageable);
+
+ @Query(nativeQuery = true,
+ value = "SELECT e.id, e.created_time, e.body, e.entity_id, e.entity_type, e.event_type, e.event_uid, e.tenant_id, ts FROM " +
+ "(SELECT *, e.body\\:\\:jsonb as json_body FROM event e WHERE " +
+ "e.tenant_id = :tenantId " +
+ "AND e.entity_type = :entityType " +
+ "AND e.entity_id = :entityId " +
+ "AND e.event_type = 'ERROR' " +
+ "AND e.created_time >= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:server IS NULL OR lower(json_body->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:method IS NULL OR lower(json_body->>'method') LIKE concat('%', lower(:method\\:\\:varchar), '%')) " +
+ "AND (:error IS NULL OR lower(json_body->>'error') LIKE concat('%', lower(:error\\:\\:varchar), '%'))",
+ countQuery = "SELECT count(*) FROM " +
+ "(SELECT *, e.body\\:\\:jsonb as json_body FROM event e WHERE " +
+ "e.tenant_id = :tenantId " +
+ "AND e.entity_type = :entityType " +
+ "AND e.entity_id = :entityId " +
+ "AND e.event_type = 'ERROR' " +
+ "AND e.created_time >= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:server IS NULL OR lower(json_body->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:method IS NULL OR lower(json_body->>'method') LIKE concat('%', lower(:method\\:\\:varchar), '%')) " +
+ "AND (:error IS NULL OR lower(json_body->>'error') LIKE concat('%', lower(:error\\:\\:varchar), '%'))")
+ Page findErrorEvents(@Param("tenantId") UUID tenantId,
+ @Param("entityId") UUID entityId,
+ @Param("entityType") String entityType,
+ @Param("startTime") Long startTime,
+ @Param("endTime") Long endTIme,
+ @Param("server") String server,
+ @Param("method") String method,
+ @Param("error") String error,
+ Pageable pageable);
+
+ @Query(nativeQuery = true,
+ value = "SELECT e.id, e.created_time, e.body, e.entity_id, e.entity_type, e.event_type, e.event_uid, e.tenant_id, ts FROM " +
+ "(SELECT *, e.body\\:\\:jsonb as json_body FROM event e WHERE " +
+ "e.tenant_id = :tenantId " +
+ "AND e.entity_type = :entityType " +
+ "AND e.entity_id = :entityId " +
+ "AND e.event_type = 'LC_EVENT' " +
+ "AND e.created_time >= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:server IS NULL OR lower(json_body->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:event IS NULL OR lower(json_body->>'event') LIKE concat('%', lower(:event\\:\\:varchar), '%')) " +
+ "AND ((:statusFilterEnabled = FALSE) OR lower(json_body->>'success')\\:\\:boolean = :statusFilter) " +
+ "AND (:error IS NULL OR lower(json_body->>'error') LIKE concat('%', lower(:error\\:\\:varchar), '%'))"
+ ,
+ countQuery = "SELECT count(*) FROM " +
+ "(SELECT *, e.body\\:\\:jsonb as json_body FROM event e WHERE " +
+ "e.tenant_id = :tenantId " +
+ "AND e.entity_type = :entityType " +
+ "AND e.entity_id = :entityId " +
+ "AND e.event_type = 'LC_EVENT' " +
+ "AND e.created_time >= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:server IS NULL OR lower(json_body->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:event IS NULL OR lower(json_body->>'event') LIKE concat('%', lower(:event\\:\\:varchar), '%')) " +
+ "AND ((:statusFilterEnabled = FALSE) OR lower(json_body->>'success')\\:\\:boolean = :statusFilter) " +
+ "AND (:error IS NULL OR lower(json_body->>'error') LIKE concat('%', lower(:error\\:\\:varchar), '%'))"
+ )
+ Page findLifeCycleEvents(@Param("tenantId") UUID tenantId,
+ @Param("entityId") UUID entityId,
+ @Param("entityType") String entityType,
+ @Param("startTime") Long startTime,
+ @Param("endTime") Long endTIme,
+ @Param("server") String server,
+ @Param("event") String event,
+ @Param("statusFilterEnabled") boolean statusFilterEnabled,
+ @Param("statusFilter") boolean statusFilter,
+ @Param("error") String error,
+ Pageable pageable);
+
+ @Query(nativeQuery = true,
+ value = "SELECT e.id, e.created_time, e.body, e.entity_id, e.entity_type, e.event_type, e.event_uid, e.tenant_id, ts FROM " +
+ "(SELECT *, e.body\\:\\:jsonb as json_body FROM event e WHERE " +
+ "e.tenant_id = :tenantId " +
+ "AND e.entity_type = :entityType " +
+ "AND e.entity_id = :entityId " +
+ "AND e.event_type = 'STATS' " +
+ "AND e.created_time >= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:server IS NULL OR lower(e.body\\:\\:json->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:messagesProcessed = 0 OR (json_body->>'messagesProcessed')\\:\\:integer >= :messagesProcessed) " +
+ "AND (:errorsOccurred = 0 OR (json_body->>'errorsOccurred')\\:\\:integer >= :errorsOccurred) ",
+ countQuery = "SELECT count(*) FROM " +
+ "(SELECT *, e.body\\:\\:jsonb as json_body FROM event e WHERE " +
+ "e.tenant_id = :tenantId " +
+ "AND e.entity_type = :entityType " +
+ "AND e.entity_id = :entityId " +
+ "AND e.event_type = 'LC_EVENT' " +
+ "AND e.created_time >= :startTime AND (:endTime = 0 OR e.created_time <= :endTime) " +
+ ") AS e WHERE " +
+ "(:server IS NULL OR lower(e.body\\:\\:json->>'server') LIKE concat('%', lower(:server\\:\\:varchar), '%')) " +
+ "AND (:messagesProcessed = 0 OR (json_body->>'messagesProcessed')\\:\\:integer >= :messagesProcessed) " +
+ "AND (:errorsOccurred = 0 OR (json_body->>'errorsOccurred')\\:\\:integer >= :errorsOccurred) ")
+ Page findStatisticsEvents(@Param("tenantId") UUID tenantId,
+ @Param("entityId") UUID entityId,
+ @Param("entityType") String entityType,
+ @Param("startTime") Long startTime,
+ @Param("endTime") Long endTIme,
+ @Param("server") String server,
+ @Param("messagesProcessed") Integer messagesProcessed,
+ @Param("errorsOccurred") Integer errorsOccurred,
+ Pageable pageable);
+
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
index 06f54042d7..44848ec515 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
@@ -24,6 +24,12 @@ import org.springframework.data.domain.PageRequest;
import org.springframework.data.repository.CrudRepository;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.Event;
+import org.thingsboard.server.common.data.event.DebugEvent;
+import org.thingsboard.server.common.data.event.ErrorEventFilter;
+import org.thingsboard.server.common.data.event.EventFilter;
+import org.thingsboard.server.common.data.event.EventType;
+import org.thingsboard.server.common.data.event.LifeCycleEventFilter;
+import org.thingsboard.server.common.data.event.StatisticsEventFilter;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.EventId;
import org.thingsboard.server.common.data.id.TenantId;
@@ -147,6 +153,98 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen
DaoUtil.toPageable(pageLink)));
}
+ @Override
+ public PageData findEventByFilter(UUID tenantId, EntityId entityId, EventFilter eventFilter, TimePageLink pageLink) {
+ if (eventFilter.hasFilterForJsonBody()) {
+ switch (eventFilter.getEventType()) {
+ case DEBUG_RULE_NODE:
+ case DEBUG_RULE_CHAIN:
+ return findEventByFilter(tenantId, entityId, (DebugEvent) eventFilter, pageLink);
+ case LC_EVENT:
+ return findEventByFilter(tenantId, entityId, (LifeCycleEventFilter) eventFilter, pageLink);
+ case ERROR:
+ return findEventByFilter(tenantId, entityId, (ErrorEventFilter) eventFilter, pageLink);
+ case STATS:
+ return findEventByFilter(tenantId, entityId, (StatisticsEventFilter) eventFilter, pageLink);
+ default:
+ throw new RuntimeException("Not supported event type: " + eventFilter.getEventType());
+ }
+ } else {
+ return findEvents(tenantId, entityId, eventFilter.getEventType().name(), pageLink);
+ }
+ }
+
+ private PageData findEventByFilter(UUID tenantId, EntityId entityId, DebugEvent eventFilter, TimePageLink pageLink) {
+ return DaoUtil.toPageData(
+ eventRepository.findDebugRuleNodeEvents(
+ tenantId,
+ entityId.getId(),
+ entityId.getEntityType().name(),
+ eventFilter.getEventType().name(),
+ notNull(pageLink.getStartTime()),
+ notNull(pageLink.getEndTime()),
+ eventFilter.getMsgDirectionType(),
+ eventFilter.getServer(),
+ eventFilter.getEntityName(),
+ eventFilter.getRelationType(),
+ eventFilter.getEntityId(),
+ eventFilter.getMsgType(),
+ eventFilter.isError(),
+ eventFilter.getError(),
+ eventFilter.getDataSearch(),
+ eventFilter.getMetadataSearch(),
+ DaoUtil.toPageable(pageLink)));
+ }
+
+ private PageData findEventByFilter(UUID tenantId, EntityId entityId, ErrorEventFilter eventFilter, TimePageLink pageLink) {
+ return DaoUtil.toPageData(
+ eventRepository.findErrorEvents(
+ tenantId,
+ entityId.getId(),
+ entityId.getEntityType().name(),
+ notNull(pageLink.getStartTime()),
+ notNull(pageLink.getEndTime()),
+ eventFilter.getServer(),
+ eventFilter.getMethod(),
+ eventFilter.getError(),
+ DaoUtil.toPageable(pageLink))
+ );
+ }
+
+ private PageData findEventByFilter(UUID tenantId, EntityId entityId, LifeCycleEventFilter eventFilter, TimePageLink pageLink) {
+ boolean statusFilterEnabled = !StringUtils.isEmpty(eventFilter.getStatus());
+ boolean statusFilter = statusFilterEnabled && eventFilter.getStatus().equalsIgnoreCase("Success");
+ return DaoUtil.toPageData(
+ eventRepository.findLifeCycleEvents(
+ tenantId,
+ entityId.getId(),
+ entityId.getEntityType().name(),
+ notNull(pageLink.getStartTime()),
+ notNull(pageLink.getEndTime()),
+ eventFilter.getServer(),
+ eventFilter.getEvent(),
+ statusFilterEnabled,
+ statusFilter,
+ eventFilter.getError(),
+ DaoUtil.toPageable(pageLink))
+ );
+ }
+
+ private PageData findEventByFilter(UUID tenantId, EntityId entityId, StatisticsEventFilter eventFilter, TimePageLink pageLink) {
+ return DaoUtil.toPageData(
+ eventRepository.findStatisticsEvents(
+ tenantId,
+ entityId.getId(),
+ entityId.getEntityType().name(),
+ notNull(pageLink.getStartTime()),
+ notNull(pageLink.getEndTime()),
+ eventFilter.getServer(),
+ notNull(eventFilter.getMessagesProcessed()),
+ notNull(eventFilter.getErrorsOccurred()),
+ DaoUtil.toPageable(pageLink))
+ );
+ }
+
@Override
public List findLatestEvents(UUID tenantId, EntityId entityId, String eventType, int limit) {
List latest = eventRepository.findLatestByTenantIdAndEntityTypeAndEntityIdAndEventType(
@@ -177,4 +275,12 @@ public class JpaBaseEventDao extends JpaAbstractDao implemen
return Optional.of(DaoUtil.getData(eventInsertRepository.saveOrUpdate(entity)));
}
+ private long notNull(Long value) {
+ return value != null ? value : 0;
+ }
+
+ private int notNull(Integer value) {
+ return value != null ? value : 0;
+ }
+
}
diff --git a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java
index 9a08c65fef..b514b43cf3 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/tenant/TenantServiceImpl.java
@@ -35,7 +35,7 @@ import org.thingsboard.server.dao.device.DeviceService;
import org.thingsboard.server.dao.entity.AbstractEntityService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.exception.DataValidationException;
-import org.thingsboard.server.dao.resource.TbResourceService;
+import org.thingsboard.server.dao.resource.ResourceService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.service.DataValidator;
import org.thingsboard.server.dao.service.PaginatedRemover;
@@ -90,7 +90,7 @@ public class TenantServiceImpl extends AbstractEntityService implements TenantSe
private RuleChainService ruleChainService;
@Autowired
- private TbResourceService resourceService;
+ private ResourceService resourceService;
@Override
public Tenant findTenantById(TenantId tenantId) {
diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
index 7d09578dd1..240d5a0b88 100644
--- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
+++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
@@ -62,6 +62,8 @@ import java.util.Collections;
import java.util.List;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.literal;
@@ -107,6 +109,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private PreparedStatement[] fetchStmtsDesc;
private PreparedStatement deleteStmt;
private PreparedStatement deletePartitionStmt;
+ private final Lock stmtCreationLock = new ReentrantLock();
private boolean isInstall() {
return environment.acceptsProfiles(Profiles.of("install"));
@@ -545,13 +548,20 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private PreparedStatement getDeleteStmt() {
if (deleteStmt == null) {
- deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF +
- " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM
- + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM
- + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM
- + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM
- + "AND " + ModelConstants.TS_COLUMN + " >= ? "
- + "AND " + ModelConstants.TS_COLUMN + " < ?");
+ stmtCreationLock.lock();
+ try {
+ if (deleteStmt == null) {
+ deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF +
+ " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM
+ + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM
+ + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM
+ + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM
+ + "AND " + ModelConstants.TS_COLUMN + " >= ? "
+ + "AND " + ModelConstants.TS_COLUMN + " < ?");
+ }
+ } finally {
+ stmtCreationLock.unlock();
+ }
}
return deleteStmt;
}
@@ -585,27 +595,41 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private PreparedStatement getDeletePartitionStmt() {
if (deletePartitionStmt == null) {
- deletePartitionStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_PARTITIONS_CF +
- " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM
- + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM
- + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM
- + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM);
+ stmtCreationLock.lock();
+ try {
+ if (deletePartitionStmt == null) {
+ deletePartitionStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_PARTITIONS_CF +
+ " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM
+ + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM
+ + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM
+ + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM);
+ }
+ } finally {
+ stmtCreationLock.unlock();
+ }
}
return deletePartitionStmt;
}
private PreparedStatement getSaveStmt(DataType dataType) {
if (saveStmts == null) {
- saveStmts = new PreparedStatement[DataType.values().length];
- for (DataType type : DataType.values()) {
- saveStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF +
- "(" + ModelConstants.ENTITY_TYPE_COLUMN +
- "," + ModelConstants.ENTITY_ID_COLUMN +
- "," + ModelConstants.KEY_COLUMN +
- "," + ModelConstants.PARTITION_COLUMN +
- "," + ModelConstants.TS_COLUMN +
- "," + getColumnName(type) + ")" +
- " VALUES(?, ?, ?, ?, ?, ?)");
+ stmtCreationLock.lock();
+ try {
+ if (saveStmts == null) {
+ saveStmts = new PreparedStatement[DataType.values().length];
+ for (DataType type : DataType.values()) {
+ saveStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF +
+ "(" + ModelConstants.ENTITY_TYPE_COLUMN +
+ "," + ModelConstants.ENTITY_ID_COLUMN +
+ "," + ModelConstants.KEY_COLUMN +
+ "," + ModelConstants.PARTITION_COLUMN +
+ "," + ModelConstants.TS_COLUMN +
+ "," + getColumnName(type) + ")" +
+ " VALUES(?, ?, ?, ?, ?, ?)");
+ }
+ }
+ } finally {
+ stmtCreationLock.unlock();
}
}
return saveStmts[dataType.ordinal()];
@@ -613,16 +637,23 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private PreparedStatement getSaveTtlStmt(DataType dataType) {
if (saveTtlStmts == null) {
- saveTtlStmts = new PreparedStatement[DataType.values().length];
- for (DataType type : DataType.values()) {
- saveTtlStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF +
- "(" + ModelConstants.ENTITY_TYPE_COLUMN +
- "," + ModelConstants.ENTITY_ID_COLUMN +
- "," + ModelConstants.KEY_COLUMN +
- "," + ModelConstants.PARTITION_COLUMN +
- "," + ModelConstants.TS_COLUMN +
- "," + getColumnName(type) + ")" +
- " VALUES(?, ?, ?, ?, ?, ?) USING TTL ?");
+ stmtCreationLock.lock();
+ try {
+ if (saveTtlStmts == null) {
+ saveTtlStmts = new PreparedStatement[DataType.values().length];
+ for (DataType type : DataType.values()) {
+ saveTtlStmts[type.ordinal()] = prepare(INSERT_INTO + ModelConstants.TS_KV_CF +
+ "(" + ModelConstants.ENTITY_TYPE_COLUMN +
+ "," + ModelConstants.ENTITY_ID_COLUMN +
+ "," + ModelConstants.KEY_COLUMN +
+ "," + ModelConstants.PARTITION_COLUMN +
+ "," + ModelConstants.TS_COLUMN +
+ "," + getColumnName(type) + ")" +
+ " VALUES(?, ?, ?, ?, ?, ?) USING TTL ?");
+ }
+ }
+ } finally {
+ stmtCreationLock.unlock();
}
}
return saveTtlStmts[dataType.ordinal()];
@@ -630,24 +661,38 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
private PreparedStatement getPartitionInsertStmt() {
if (partitionInsertStmt == null) {
- partitionInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF +
- "(" + ModelConstants.ENTITY_TYPE_COLUMN +
- "," + ModelConstants.ENTITY_ID_COLUMN +
- "," + ModelConstants.PARTITION_COLUMN +
- "," + ModelConstants.KEY_COLUMN + ")" +
- " VALUES(?, ?, ?, ?)");
+ stmtCreationLock.lock();
+ try {
+ if (partitionInsertStmt == null) {
+ partitionInsertStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF +
+ "(" + ModelConstants.ENTITY_TYPE_COLUMN +
+ "," + ModelConstants.ENTITY_ID_COLUMN +
+ "," + ModelConstants.PARTITION_COLUMN +
+ "," + ModelConstants.KEY_COLUMN + ")" +
+ " VALUES(?, ?, ?, ?)");
+ }
+ } finally {
+ stmtCreationLock.unlock();
+ }
}
return partitionInsertStmt;
}
private PreparedStatement getPartitionInsertTtlStmt() {
if (partitionInsertTtlStmt == null) {
- partitionInsertTtlStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF +
- "(" + ModelConstants.ENTITY_TYPE_COLUMN +
- "," + ModelConstants.ENTITY_ID_COLUMN +
- "," + ModelConstants.PARTITION_COLUMN +
- "," + ModelConstants.KEY_COLUMN + ")" +
- " VALUES(?, ?, ?, ?) USING TTL ?");
+ stmtCreationLock.lock();
+ try {
+ if (partitionInsertTtlStmt == null) {
+ partitionInsertTtlStmt = prepare(INSERT_INTO + ModelConstants.TS_KV_PARTITIONS_CF +
+ "(" + ModelConstants.ENTITY_TYPE_COLUMN +
+ "," + ModelConstants.ENTITY_ID_COLUMN +
+ "," + ModelConstants.PARTITION_COLUMN +
+ "," + ModelConstants.KEY_COLUMN + ")" +
+ " VALUES(?, ?, ?, ?) USING TTL ?");
+ }
+ } finally {
+ stmtCreationLock.unlock();
+ }
}
return partitionInsertTtlStmt;
}
@@ -713,12 +758,26 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
switch (orderBy) {
case ASC_ORDER:
if (fetchStmtsAsc == null) {
- fetchStmtsAsc = initFetchStmt(orderBy);
+ stmtCreationLock.lock();
+ try {
+ if (fetchStmtsAsc == null) {
+ fetchStmtsAsc = initFetchStmt(orderBy);
+ }
+ } finally {
+ stmtCreationLock.unlock();
+ }
}
return fetchStmtsAsc[aggType.ordinal()];
case DESC_ORDER:
if (fetchStmtsDesc == null) {
- fetchStmtsDesc = initFetchStmt(orderBy);
+ stmtCreationLock.lock();
+ try {
+ if (fetchStmtsDesc == null) {
+ fetchStmtsDesc = initFetchStmt(orderBy);
+ }
+ } finally {
+ stmtCreationLock.unlock();
+ }
}
return fetchStmtsDesc[aggType.ordinal()];
default:
diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java
index cfd232b330..0761f28d33 100644
--- a/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java
+++ b/dao/src/test/java/org/thingsboard/server/dao/service/AbstractServiceTest.java
@@ -54,7 +54,7 @@ import org.thingsboard.server.dao.entity.EntityService;
import org.thingsboard.server.dao.entityview.EntityViewService;
import org.thingsboard.server.dao.event.EventService;
import org.thingsboard.server.dao.relation.RelationService;
-import org.thingsboard.server.dao.resource.TbResourceService;
+import org.thingsboard.server.dao.resource.ResourceService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.settings.AdminSettingsService;
import org.thingsboard.server.dao.tenant.TenantProfileService;
@@ -152,9 +152,9 @@ public abstract class AbstractServiceTest {
protected DeviceProfileService deviceProfileService;
@Autowired
- protected TbResourceService resourceService;
+ protected ResourceService resourceService;
- class IdComparator implements Comparator {
+ public class IdComparator implements Comparator {
@Override
public int compare(D o1, D o2) {
return o1.getId().getId().compareTo(o2.getId().getId());
diff --git a/pom.xml b/pom.xml
index c26964e8e9..d954128e96 100755
--- a/pom.xml
+++ b/pom.xml
@@ -99,7 +99,7 @@
org/thingsboard/server/extensions/core/plugin/telemetry/gen/**/*
5.0.2
- 0.1.31
+ 0.1.16
2.6.0
4.1.1
2.57
diff --git a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java
index de7ba29141..815897faaa 100644
--- a/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java
+++ b/rest-client/src/main/java/org/thingsboard/rest/client/RestClient.java
@@ -1338,7 +1338,7 @@ public class RestClient implements ClientHttpRequestInterceptor, Closeable {
HttpEntity.EMPTY, DeviceProfile.class, deviceProfileId).getBody();
}
- public PageData getTenantDevices(PageLink pageLink) {
+ public PageData getDeviceProfiles(PageLink pageLink) {
Map params = new HashMap<>();
addPageLinkToParam(params, pageLink);
return restTemplate.exchange(
diff --git a/ui-ngx/src/app/core/http/event.service.ts b/ui-ngx/src/app/core/http/event.service.ts
index 9bf1b4e86f..fd740d73a7 100644
--- a/ui-ngx/src/app/core/http/event.service.ts
+++ b/ui-ngx/src/app/core/http/event.service.ts
@@ -21,7 +21,7 @@ import { HttpClient } from '@angular/common/http';
import { TimePageLink } from '@shared/models/page/page-link';
import { PageData } from '@shared/models/page/page-data';
import { EntityId } from '@shared/models/id/entity-id';
-import { DebugEventType, Event, EventType } from '@shared/models/event.models';
+import { DebugEventType, Event, EventType, FilterEventBody } from '@shared/models/event.models';
@Injectable({
providedIn: 'root'
@@ -39,4 +39,10 @@ export class EventService {
defaultHttpOptionsFromConfig(config));
}
+ public getFilterEvents(entityId: EntityId, eventType: EventType | DebugEventType, tenantId: string,
+ filters: FilterEventBody, pageLink: TimePageLink, config?: RequestConfig): Observable> {
+ return this.http.post>(`/api/events/${entityId.entityType}/${entityId.id}` +
+ `${pageLink.toQuery()}&tenantId=${tenantId}`, {...filters, eventType}, defaultHttpOptionsFromConfig(config));
+ }
+
}
diff --git a/ui-ngx/src/app/core/http/resource.service.ts b/ui-ngx/src/app/core/http/resource.service.ts
index d8bc18291b..2fef72fedb 100644
--- a/ui-ngx/src/app/core/http/resource.service.ts
+++ b/ui-ngx/src/app/core/http/resource.service.ts
@@ -18,10 +18,10 @@ import { Injectable } from '@angular/core';
import { HttpClient } from '@angular/common/http';
import { PageLink } from '@shared/models/page/page-link';
import { defaultHttpOptionsFromConfig, RequestConfig } from '@core/http/http-utils';
-import { Observable } from 'rxjs';
+import { forkJoin, Observable, of } from 'rxjs';
import { PageData } from '@shared/models/page/page-data';
import { Resource, ResourceInfo } from '@shared/models/resource.models';
-import { map } from 'rxjs/operators';
+import { catchError, map, mergeMap } from 'rxjs/operators';
@Injectable({
providedIn: 'root'
@@ -43,14 +43,17 @@ export class ResourceService {
}
public downloadResource(resourceId: string): Observable {
- return this.http.get(`/api/resource/${resourceId}/download`, { responseType: 'arraybuffer', observe: 'response' }).pipe(
+ return this.http.get(`/api/resource/${resourceId}/download`, {
+ responseType: 'arraybuffer',
+ observe: 'response'
+ }).pipe(
map((response) => {
const headers = response.headers;
const filename = headers.get('x-filename');
const contentType = headers.get('content-type');
const linkElement = document.createElement('a');
try {
- const blob = new Blob([response.body], { type: contentType });
+ const blob = new Blob([response.body], {type: contentType});
const url = URL.createObjectURL(blob);
linkElement.setAttribute('href', url);
linkElement.setAttribute('download', filename);
@@ -70,6 +73,25 @@ export class ResourceService {
);
}
+ public saveResources(resources: Resource[], config?: RequestConfig): Observable {
+ let partSize = 100;
+ partSize = resources.length > partSize ? partSize : resources.length;
+ const resourceObservables: Observable[] = [];
+ for (let i = 0; i < partSize; i++) {
+ resourceObservables.push(this.saveResource(resources[i], config).pipe(catchError(() => of({} as Resource))));
+ }
+ return forkJoin(resourceObservables).pipe(
+ mergeMap((resource) => {
+ resources.splice(0, partSize);
+ if (resources.length) {
+ return this.saveResources(resources, config);
+ } else {
+ return of(resource);
+ }
+ })
+ );
+ }
+
public saveResource(resource: Resource, config?: RequestConfig): Observable {
return this.http.post('/api/resource', resource, defaultHttpOptionsFromConfig(config));
}
diff --git a/ui-ngx/src/app/core/services/time.service.ts b/ui-ngx/src/app/core/services/time.service.ts
index 47c8d8a456..23fc215573 100644
--- a/ui-ngx/src/app/core/services/time.service.ts
+++ b/ui-ngx/src/app/core/services/time.service.ts
@@ -31,7 +31,7 @@ import { isDefined } from '@core/utils';
export interface TimeInterval {
name: string;
- translateParams: {[key: string]: any};
+ translateParams: { [key: string]: any };
value: number;
}
@@ -56,14 +56,14 @@ export class TimeService {
public loadMaxDatapointsLimit(): Observable {
return this.http.get('/api/dashboard/maxDatapointsLimit',
defaultHttpOptions(true)).pipe(
- map( (limit) => {
- this.maxDatapointsLimit = limit;
- if (!this.maxDatapointsLimit || this.maxDatapointsLimit <= MIN_LIMIT) {
- this.maxDatapointsLimit = MIN_LIMIT + 1;
- }
- return this.maxDatapointsLimit;
- })
- );
+ map((limit) => {
+ this.maxDatapointsLimit = limit;
+ if (!this.maxDatapointsLimit || this.maxDatapointsLimit <= MIN_LIMIT) {
+ this.maxDatapointsLimit = MIN_LIMIT + 1;
+ }
+ return this.maxDatapointsLimit;
+ })
+ );
}
public matchesExistingInterval(min: number, max: number, intervalMs: number): boolean {
@@ -79,7 +79,7 @@ export class TimeService {
public boundMinInterval(min: number): number {
if (isDefined(min)) {
- min = Math.floor(min / 1000) * 1000;
+ min = Math.ceil(min / 1000) * 1000;
}
return this.toBound(min, MIN_INTERVAL, MAX_INTERVAL, MIN_INTERVAL);
}
diff --git a/ui-ngx/src/app/core/utils.ts b/ui-ngx/src/app/core/utils.ts
index b9c2697d0a..b2d7b8b2a7 100644
--- a/ui-ngx/src/app/core/utils.ts
+++ b/ui-ngx/src/app/core/utils.ts
@@ -127,6 +127,10 @@ export function isEmpty(obj: any): boolean {
return true;
}
+export function isLiteralObject(value: any) {
+ return (!!value) && (value.constructor === Object);
+}
+
export function formatValue(value: any, dec?: number, units?: string, showZeroDecimals?: boolean): string | undefined {
if (isDefinedAndNotNull(value) && isNumeric(value) &&
(isDefinedAndNotNull(dec) || isDefinedAndNotNull(units) || Number(value).toString() === value)) {
@@ -287,7 +291,7 @@ export function deepClone(target: T, ignoreFields?: string[]): T {
return cp.map((n: any) => deepClone(n)) as any;
}
if (typeof target === 'object' && target !== {}) {
- const cp = { ...(target as { [key: string]: any }) } as { [key: string]: any };
+ const cp = {...(target as { [key: string]: any })} as { [key: string]: any };
Object.keys(cp).forEach(k => {
if (!ignoreFields || ignoreFields.indexOf(k) === -1) {
cp[k] = deepClone(cp[k]);
diff --git a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html
index 762b7149a2..7344a14310 100644
--- a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html
+++ b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.html
@@ -56,6 +56,8 @@
[syncStateWithQueryParam]="syncStateWithQueryParam"
[states]="dashboardConfiguration.states">
+
diff --git a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.scss b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.scss
index c7d2367975..6ace69d312 100644
--- a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.scss
+++ b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.scss
@@ -54,6 +54,11 @@ div.tb-dashboard-page {
z-index: 13;
pointer-events: none;
+ .dashboard_logo{
+ height: 75%;
+ margin-right: 16px;
+ }
+
&.tb-dashboard-toolbar-opened {
right: 0;
// transition: right .3s cubic-bezier(.55, 0, .55, .2);
diff --git a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts
index 7e760e3337..738304cf64 100644
--- a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts
+++ b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-page.component.ts
@@ -90,7 +90,6 @@ import {
} from '@home/components/alias/entity-aliases-dialog.component';
import { EntityAliases } from '@app/shared/models/alias.models';
import { EditWidgetComponent } from '@home/components/dashboard-page/edit-widget.component';
-import { WidgetsBundle } from '@shared/models/widgets-bundle.model';
import {
AddWidgetDialogComponent,
AddWidgetDialogData
@@ -118,8 +117,7 @@ import { ComponentPortal } from '@angular/cdk/portal';
import {
DISPLAY_WIDGET_TYPES_PANEL_DATA,
DisplayWidgetTypesPanelComponent,
- DisplayWidgetTypesPanelData,
- WidgetTypes
+ DisplayWidgetTypesPanelData
} from '@home/components/dashboard-page/widget-types-panel.component';
import { DashboardWidgetSelectComponent } from '@home/components/dashboard-page/dashboard-widget-select.component';
import {AliasEntityType, EntityType} from "@shared/models/entity-type.models";
@@ -193,6 +191,7 @@ export class DashboardPageComponent extends PageComponent implements IDashboardC
addingLayoutCtx: DashboardPageLayoutContext;
+ logo = 'assets/logo_title_white.svg';
dashboardCtx: DashboardContext = {
instanceId: this.utils.guid(),
@@ -486,6 +485,19 @@ export class DashboardPageComponent extends PageComponent implements IDashboardC
}
}
+ public showDashboardLogo(): boolean {
+ if (this.dashboard.configuration.settings &&
+ isDefined(this.dashboard.configuration.settings.showDashboardLogo)) {
+ return this.dashboard.configuration.settings.showDashboardLogo;
+ } else {
+ return false;
+ }
+ }
+
+ public get dashboardLogo(): string {
+ return this.dashboard.configuration.settings.dashboardLogoUrl || this.logo;
+ }
+
public showRightLayoutSwitch(): boolean {
return this.isMobile && this.layouts.right.show;
}
@@ -607,7 +619,7 @@ export class DashboardPageComponent extends PageComponent implements IDashboardC
panelClass: ['tb-dialog', 'tb-fullscreen-dialog'],
data: {
settings: deepClone(this.dashboard.configuration.settings),
- gridSettings
+ gridSettings,
}
}).afterClosed().subscribe((data) => {
if (data) {
@@ -1189,13 +1201,16 @@ export class DashboardPageComponent extends PageComponent implements IDashboardC
overlayRef.dispose();
});
+ const filterWidgetTypes = this.dashboardWidgetSelectComponent.filterWidgetTypes;
+ const widgetTypesList = Array.from(this.dashboardWidgetSelectComponent.widgetTypes.values()).map(type => {
+ return {type, display: filterWidgetTypes === null ? true : filterWidgetTypes.includes(type)};
+ });
+
const providers: StaticProvider[] = [
{
provide: DISPLAY_WIDGET_TYPES_PANEL_DATA,
useValue: {
- types: Array.from(this.dashboardWidgetSelectComponent.widgetTypes.values()).map(type => {
- return {type, display: true};
- }),
+ types: widgetTypesList,
typesUpdated: (newTypes) => {
this.filterWidgetTypes = newTypes.filter(type => type.display).map(type => type.type);
}
diff --git a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-settings-dialog.component.html b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-settings-dialog.component.html
index ede8052c68..3a20d146dd 100644
--- a/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-settings-dialog.component.html
+++ b/ui-ngx/src/app/modules/home/components/dashboard-page/dashboard-settings-dialog.component.html
@@ -52,7 +52,8 @@
formControlName="titleColor">
-
+
{{ 'dashboard.display-dashboards-selection' | translate }}
@@ -69,6 +70,13 @@
{{ 'dashboard.display-dashboard-export' | translate }}
+
+ {{ 'dashboard.display-dashboard-toolbar-logo' | translate }}
+
+
+
{
+ return this.filterWidgetTypes$.value;
+ }
+
@Output()
widgetSelected: EventEmitter = new EventEmitter();
diff --git a/ui-ngx/src/app/modules/home/components/dashboard/dashboard.component.html b/ui-ngx/src/app/modules/home/components/dashboard/dashboard.component.html
index d893ec848f..c050b524cf 100644
--- a/ui-ngx/src/app/modules/home/components/dashboard/dashboard.component.html
+++ b/ui-ngx/src/app/modules/home/components/dashboard/dashboard.component.html
@@ -80,7 +80,7 @@
(mousedown)="widgetMouseDown($event, widget)"
(click)="widgetClicked($event, widget)"
(contextmenu)="openWidgetContextMenu($event, widget)">
-