From 1869296ff240df84704d522791b694e564dac66a Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Thu, 17 Apr 2025 12:02:45 +0300 Subject: [PATCH 01/28] added ability to preserve last update ts for calculated value --- .../controller/CalculatedFieldController.java | 19 ++++++++- .../ctx/state/BaseCalculatedFieldState.java | 28 +++++++++++-- .../cf/ctx/state/CalculatedFieldCtx.java | 3 ++ .../cf/ctx/state/CalculatedFieldState.java | 2 + .../ctx/state/ScriptCalculatedFieldState.java | 3 +- .../ctx/state/SimpleCalculatedFieldState.java | 39 +++++++++++++------ .../SimpleCalculatedFieldConfiguration.java | 2 + .../script/api/tbel/TbelCfCtx.java | 5 ++- 8 files changed, 82 insertions(+), 19 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java index f899d0f480..7e4e041600 100644 --- a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java +++ b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java @@ -36,6 +36,8 @@ import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.script.api.tbel.TbelCfArg; import org.thingsboard.script.api.tbel.TbelCfCtx; import org.thingsboard.script.api.tbel.TbelCfSingleValueArg; +import org.thingsboard.script.api.tbel.TbelCfTsDoubleVal; +import org.thingsboard.script.api.tbel.TbelCfTsRollingArg; import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EventInfo; @@ -242,9 +244,8 @@ public class CalculatedFieldController extends BaseController { ctxAndArgNames.toArray(String[]::new) ); - Object[] args = new Object[ctxAndArgNames.size()]; - args[0] = new TbelCfCtx(arguments); + args[0] = new TbelCfCtx(arguments, getLastUpdateTimestamp(arguments)); for (int i = 1; i < ctxAndArgNames.size(); i++) { var arg = arguments.get(ctxAndArgNames.get(i)); if (arg instanceof TbelCfSingleValueArg svArg) { @@ -267,6 +268,20 @@ public class CalculatedFieldController extends BaseController { return result; } + private long getLastUpdateTimestamp(Map arguments) { + long lastUpdateTimestamp = -1; + for (TbelCfArg entry : arguments.values()) { + if (entry instanceof TbelCfSingleValueArg singleValueArg) { + long ts = singleValueArg.getTs(); + lastUpdateTimestamp = Math.max(lastUpdateTimestamp, ts); + } else if (entry instanceof TbelCfTsRollingArg tsRollingArg) { + long maxTs = tsRollingArg.getValues().stream().mapToLong(TbelCfTsDoubleVal::getTs).max().orElse(-1); + lastUpdateTimestamp = Math.max(lastUpdateTimestamp, maxTs); + } + } + return lastUpdateTimestamp; + } + private & HasTenantId, I extends EntityId> void checkReferencedEntities(CalculatedFieldConfiguration calculatedFieldConfig, SecurityUser user) throws ThingsboardException { List referencedEntityIds = calculatedFieldConfig.getReferencedEntities(); for (EntityId referencedEntityId : referencedEntityIds) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index 80b003b3cc..84c61661ae 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -35,13 +35,15 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { protected Map arguments; protected boolean sizeExceedsLimit; + protected long lastUpdateTimestamp = -1; + public BaseCalculatedFieldState(List requiredArguments) { this.requiredArguments = requiredArguments; this.arguments = new HashMap<>(); } public BaseCalculatedFieldState() { - this(new ArrayList<>(), new HashMap<>(), false); + this(new ArrayList<>(), new HashMap<>(), false, -1); } @Override @@ -59,14 +61,21 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { checkArgumentSize(key, newEntry, ctx); ArgumentEntry existingEntry = arguments.get(key); + boolean entryUpdated; if (existingEntry == null || newEntry.isForceResetPrevious()) { validateNewEntry(newEntry); arguments.put(key, newEntry); - stateUpdated = true; + entryUpdated = true; } else { - stateUpdated = existingEntry.updateEntry(newEntry); + entryUpdated = existingEntry.updateEntry(newEntry); + } + + if (entryUpdated) { + stateUpdated = true; + updateLastUpdateTimestamp(newEntry); } + } return stateUpdated; @@ -100,4 +109,17 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { protected abstract void validateNewEntry(ArgumentEntry newEntry); + private void updateLastUpdateTimestamp(ArgumentEntry entry) { + if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { + long ts = singleValueArgumentEntry.getTs(); + this.lastUpdateTimestamp = Math.max(this.lastUpdateTimestamp, ts); + } else if (entry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { + Map.Entry lastEntry = tsRollingArgumentEntry.getTsRecords().pollLastEntry(); + if (lastEntry != null) { + long ts = lastEntry.getKey(); + this.lastUpdateTimestamp = Math.max(this.lastUpdateTimestamp, ts); + } + } + } + } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index 674ce3726f..25083ed415 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -28,6 +28,7 @@ import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; +import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -61,6 +62,7 @@ public class CalculatedFieldCtx { private final List argNames; private Output output; private String expression; + private boolean preserveLatestTs; private TbelInvokeService tbelInvokeService; private CalculatedFieldScriptEngine calculatedFieldScriptEngine; private ThreadLocal customExpression; @@ -94,6 +96,7 @@ public class CalculatedFieldCtx { this.argNames = new ArrayList<>(arguments.keySet()); this.output = configuration.getOutput(); this.expression = configuration.getExpression(); + this.preserveLatestTs = CalculatedFieldType.SIMPLE.equals(calculatedField.getType()) && ((SimpleCalculatedFieldConfiguration) configuration).isPreserveLastUpdateTs(); this.tbelInvokeService = tbelInvokeService; this.maxDataPointsPerRollingArg = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java index fc4ba513d2..6eac3358ba 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java @@ -42,6 +42,8 @@ public interface CalculatedFieldState { Map getArguments(); + long getLastUpdateTimestamp(); + void setRequiredArguments(List requiredArguments); boolean updateState(CalculatedFieldCtx ctx, Map argumentValues); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java index bf00f1b0b1..65ef40330c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java @@ -30,7 +30,6 @@ import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.service.cf.CalculatedFieldResult; import java.util.ArrayList; -import java.util.HashMap; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -67,7 +66,7 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { args.add(arg); } } - args.set(0, new TbelCfCtx(arguments)); + args.set(0, new TbelCfCtx(arguments, getLastUpdateTimestamp())); ListenableFuture resultFuture = ctx.getCalculatedFieldScriptEngine().executeJsonAsync(args.toArray()); Output output = ctx.getOutput(); return Futures.transform(resultFuture, diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 480b334ac3..2ad961f72c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -15,6 +15,8 @@ */ package org.thingsboard.server.service.cf.ctx.state; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.node.ObjectNode; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.Data; @@ -65,19 +67,34 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { double expressionResult = expr.evaluate(); Output output = ctx.getOutput(); - Object result; - Integer decimals = output.getDecimalsByDefault(); - if (decimals != null) { - if (decimals.equals(0)) { - result = TbUtils.toInt(expressionResult); - } else { - result = TbUtils.toFixed(expressionResult, decimals); - } - } else { - result = expressionResult; + Object result = formatResult(expressionResult, output.getDecimalsByDefault()); + JsonNode outputResult = createResultJson(ctx.isPreserveLatestTs(), output.getName(), result); + + return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), outputResult)); + } + + private Object formatResult(double expressionResult, Integer decimals) { + if (decimals == null) { + return expressionResult; } + return decimals.equals(0) + ? TbUtils.toInt(expressionResult) + : TbUtils.toFixed(expressionResult, decimals); + } + + private JsonNode createResultJson(boolean preserveLatestTs, String outputName, Object result) { + ObjectNode valuesNode = JacksonUtil.newObjectNode(); + valuesNode.set(outputName, JacksonUtil.valueToTree(result)); - return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), JacksonUtil.valueToTree(Map.of(output.getName(), result)))); + long lastTimestamp = getLastUpdateTimestamp(); + if (preserveLatestTs && lastTimestamp != -1) { + ObjectNode resultNode = JacksonUtil.newObjectNode(); + resultNode.put("ts", lastTimestamp); + resultNode.set("values", valuesNode); + return resultNode; + } else { + return valuesNode; + } } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java index 5c0ce71e86..c748afdac4 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java @@ -23,6 +23,8 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldType; @EqualsAndHashCode(callSuper = true) public class SimpleCalculatedFieldConfiguration extends BaseCalculatedFieldConfiguration implements CalculatedFieldConfiguration { + private boolean preserveLastUpdateTs; + @Override public CalculatedFieldType getType() { return CalculatedFieldType.SIMPLE; diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java index ce42e2cf3b..1a8610cbab 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java @@ -24,9 +24,12 @@ public class TbelCfCtx implements TbelCfObject { @Getter private final Map args; + @Getter + private final long lastTs; - public TbelCfCtx(Map args) { + public TbelCfCtx(Map args, long lastUpdateTs) { this.args = Collections.unmodifiableMap(args); + this.lastTs = lastUpdateTs != -1 ? lastUpdateTs : System.currentTimeMillis(); } @Override From 6f9e36d386cba656deb5a57114566d670a22e941 Mon Sep 17 00:00:00 2001 From: yevhenii Date: Thu, 24 Apr 2025 18:36:39 +0300 Subject: [PATCH 02/28] Update edge install instructions - updated install instructions --- .../install/centos/instructions.md | 144 ++++++++---------- .../install/docker/instructions.md | 61 +++++--- .../install/ubuntu/instructions.md | 112 +++++++------- 3 files changed, 153 insertions(+), 164 deletions(-) diff --git a/application/src/main/data/json/edge/instructions/install/centos/instructions.md b/application/src/main/data/json/edge/instructions/install/centos/instructions.md index 86c5acad56..90d6d6ddc8 100644 --- a/application/src/main/data/json/edge/instructions/install/centos/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/centos/instructions.md @@ -4,15 +4,15 @@ Here is the list of commands, that can be used to quickly install ThingsBoard Ed Before continue to installation execute the following commands in order to install necessary tools: ```bash -sudo yum install -y nano wget -sudo yum install -y https://dl.fedoraproject.org/pub/epel/epel-release-latest-7.noarch.rpm +sudo yum install -y nano wget && sudo yum install -y https://dl.fedoraproject.org/pub/epel/epel-release-latest-7.noarch.rpm +{:copy-code} ``` -#### Install Java 17 (OpenJDK) +#### Step 1. Install Java 17 (OpenJDK) ThingsBoard service is running on Java 17. Follow these instructions to install OpenJDK 17: ```bash -sudo yum install java-17-openjdk +sudo dnf install java-17-openjdk {:copy-code} ``` @@ -39,112 +39,92 @@ OpenJDK Runtime Environment (...) OpenJDK 64-Bit Server VM (build ...) ``` -#### Configure PostgreSQL +#### Step 2. Configure ThingsBoard Database +ThingsBoard Edge supports SQL and hybrid database approaches. + +### PostgresSql ThingsBoard Edge uses PostgreSQL database as a local storage. -Instructions listed below will help you to install PostgreSQL. +To install PostgreSQL, follow the instructions below. ```bash # Update your system -sudo yum update +sudo dnf update {:copy-code} ``` -**For CentOS 7:** +Install the repository RPM: +**For CentOS/RHEL 8:** ```bash -# Install the repository RPM (for CentOS 7): -sudo yum -y install https://download.postgresql.org/pub/repos/yum/reporpms/EL-7-x86_64/pgdg-redhat-repo-latest.noarch.rpm -# Install packages -sudo yum -y install epel-release yum-utils -sudo yum-config-manager --enable pgdg16 -sudo yum install postgresql16-server postgresql16 postgresql16-contrib -# Initialize your PostgreSQL DB -sudo /usr/pgsql-16/bin/postgresql-16-setup initdb -sudo systemctl start postgresql-16 -# Optional: Configure PostgreSQL to start on boot -sudo systemctl enable --now postgresql-16 - +# Install the repository RPM (For CentOS/RHEL 8): +sudo sudo dnf -y install https://download.postgresql.org/pub/repos/yum/reporpms/EL-8-x86_64/pgdg-redhat-repo-latest.noarch.rpm {:copy-code} ``` -**For CentOS 8:** +**For CentOS/RHEL 9:** ```bash -# Install the repository RPM (for CentOS 8): -sudo yum -y install https://download.postgresql.org/pub/repos/yum/reporpms/EL-8-x86_64/pgdg-redhat-repo-latest.noarch.rpm -# Install packages -sudo dnf -qy module disable postgresql -sudo dnf -y install postgresql16 postgresql16-server postgresql16-contrib -# Initialize your PostgreSQL DB -sudo /usr/pgsql-16/bin/postgresql-16-setup initdb -sudo systemctl start postgresql-16 -# Optional: Configure PostgreSQL to start on boot -sudo systemctl enable --now postgresql-16 - +# Install the repository RPM (for CentOS 9): +sudo dnf -y install https://download.postgresql.org/pub/repos/yum/reporpms/EL-9-x86_64/pgdg-redhat-repo-latest.noarch.rpm {:copy-code} ``` -Once PostgreSQL is installed you may want to create a new user or set the password for the main user. -The instructions below will help to set the password for main PostgreSQL user: +Install packages and initialize PostgreSQL. The PostgreSQL service will automatically start every time the system boots up. -```text -sudo su - postgres -psql -\password -\q +```bash +sudo dnf -qy module disable postgresql && \ +sudo dnf -y install postgresql16 postgresql16-server postgresql16-contrib && \ +sudo /usr/pgsql-16/bin/postgresql-16-setup initdb && \ +sudo systemctl enable --now postgresql-16 +{:copy-code} ``` -Then, press "Ctrl+D" to return to main user console. +Once PostgreSQL is installed, it is recommended to set the password for the PostgreSQL main user. -After configuring the password, edit the pg_hba.conf to use MD5 authentication with the postgres user. - -Edit pg_hba.conf file: +The following command will switch the current user to the PostgreSQL user and set the password directly in PostgreSQL. ```bash -sudo nano /var/lib/pgsql/16/data/pg_hba.conf +sudo -u postgres psql -c "\password" {:copy-code} ``` -Locate the following lines: - -```text -# IPv4 local connections: -host all all 127.0.0.1/32 ident -``` - -Replace `ident` with `md5`: +Then, enter and confirm the password. -```text -host all all 127.0.0.1/32 md5 -``` +Since ThingsBoard Edge uses the PostgreSQL database for local storage, configuring MD5 authentication ensures that only authenticated users or +applications can access the database, thus protecting your data. After configuring the password, +edit the pg_hba.conf file to use MD5 hashing for authentication instead of the default method (ident) for local IPv4 connections. -Finally, you should restart the PostgreSQL service to initialize the new configuration: +To replace ident with md5, run the following command: ```bash -sudo systemctl restart postgresql-16.service +sudo sed -i 's/^host\s\+all\s\+all\s\+127\.0\.0\.1\/32\s\+ident/host all all 127.0.0.1\/32 md5/' /var/lib/pgsql/16/data/pg_hba.conf {:copy-code} ``` -Connect to the database to create ThingsBoard Edge DB: +Then run the command that will restart the PostgreSQL service to apply configuration changes, connect to the database as a postgres user, +and create the ThingsBoard Edge database (tb_edge). To connect to the PostgreSQL database, enter the PostgreSQL password. ```bash -psql -U postgres -d postgres -h 127.0.0.1 -W +sudo systemctl restart postgresql-16.service && psql -U postgres -d postgres -h 127.0.0.1 -W -c "CREATE DATABASE tb_edge;" {:copy-code} ``` -Execute create database statement: +#### Step 3. Choose Queue Service -```bash -CREATE DATABASE tb_edge; -\q -{:copy-code} -``` +ThingsBoard Edge is able to use different messaging systems/brokers for storing the messages and communication between ThingsBoard services. +How to choose the right queue implementation? -#### ThingsBoard Edge service installation +In Memory queue implementation is built-in and default. It is useful for development(PoC) environments and is not suitable for production deployments or any sort of cluster deployments. + +Kafka is recommended for production deployments. This queue is used on the most of ThingsBoard production environments now. + +In Memory queue is built in and enabled by default. No additional configuration is required. + +#### Step 4. ThingsBoard Edge Service Installation Download installation package: ```bash -wget https://github.com/thingsboard/thingsboard-edge/releases/download/v${TB_EDGE_TAG}/tb-edge-${TB_EDGE_TAG}.rpm +wget wget https://github.com/thingsboard/thingsboard-edge/releases/download/v${TB_EDGE_TAG}/tb-edge-${TB_EDGE_TAG}.rpm {:copy-code} ``` @@ -155,7 +135,7 @@ sudo rpm -Uvh tb-edge-${TB_EDGE_TAG}.rpm {:copy-code} ``` -#### Configure ThingsBoard Edge +#### Step 5. Configure ThingsBoard Edge To configure ThingsBoard Edge, you can use the following command to automatically update the configuration file with specific values: ```bash @@ -169,25 +149,21 @@ EOL' {:copy-code} ``` -##### [Optional] Database Configuration -In case you changed default PostgreSQL datasource settings (**postgres**/**postgres**) please update the configuration file (**/etc/tb-edge/conf/tb-edge.conf**) with your actual values: - -```bash -sudo nano /etc/tb-edge/conf/tb-edge.conf -{:copy-code} -``` - -Please update the following lines in your configuration file. Make sure **to replace**: -- Replace 'postgres' with your actual PostgreSQL username; -- Replace 'PUT_YOUR_POSTGRESQL_PASSWORD_HERE' with your actual PostgreSQL password. +##### Configure PostgreSQL (Optional) +If you changed PostgreSQL default datasource settings, use the following command: ```bash +sudo sh -c 'cat <> /etc/tb-edge/conf/tb-edge.conf export SPRING_DATASOURCE_URL=jdbc:postgresql://localhost:5432/tb_edge export SPRING_DATASOURCE_USERNAME=postgres -export SPRING_DATASOURCE_PASSWORD=PUT_YOUR_POSTGRESQL_PASSWORD_HERE +export SPRING_DATASOURCE_PASSWORD= +EOL' {:copy-code} ``` +PUT_YOUR_POSTGRESQL_PASSWORD_HERE: Replace with your actual PostgreSQL user password. + + ##### [Optional] Update bind ports If ThingsBoard Edge is going to be running on the same machine where ThingsBoard server (cloud) is running, you'll need to update configuration parameters to avoid port collision between ThingsBoard server and ThingsBoard Edge. @@ -206,7 +182,7 @@ EOL' Make sure that ports above (18080, 11883, 15683) are not used by any other application. -#### Run installation script +#### Step 6. Run installation Script Once ThingsBoard Edge is installed and configured please execute the following install script: ```bash @@ -214,18 +190,18 @@ sudo /usr/share/tb-edge/bin/install/install.sh {:copy-code} ``` -#### Restart ThingsBoard Edge service +#### Step 7. Restart ThingsBoard Edge Service ```bash sudo service tb-edge restart {:copy-code} ``` -#### Open ThingsBoard Edge UI +#### Step 8. Open ThingsBoard Edge UI Once started, you will be able to open **ThingsBoard Edge UI** using the following link http://localhost:8080. ###### NOTE: Edge HTTP bind port update -Use next **ThingsBoard Edge UI** link **http://localhost:18080** if you updated HTTP 8080 bind port to **18080**. +If the Edge HTTP bind port was changed to 18080 during Edge installation, access the ThingsBoard Edge instance at http://localhost:18080. diff --git a/application/src/main/data/json/edge/instructions/install/docker/instructions.md b/application/src/main/data/json/edge/instructions/install/docker/instructions.md index 23f484f2ec..4457962d86 100644 --- a/application/src/main/data/json/edge/instructions/install/docker/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/docker/instructions.md @@ -4,9 +4,26 @@ Here is the list of commands, that can be used to quickly install ThingsBoard Ed Install Docker CE and Docker Compose. -#### Running ThingsBoard Edge as docker service +#### Step 1. Running ThingsBoard Edge -Create docker compose file for ThingsBoard Edge service: +Here you can find ThingsBoard Edge docker image: + + thingsboard/tb-edge + +#### Step 2. Choose Queue and/or Database Services + +ThingsBoard Edge is able to use different messaging systems/brokers for storing the messages and communication between ThingsBoard services. +How to choose the right queue implementation? + +In Memory queue implementation is built-in and default. It is useful for development(PoC) environments and is not suitable for production deployments or any sort of cluster deployments. + +Kafka is recommended for production deployments. This queue is used on the most of ThingsBoard production environments now. + +Hybrid implementation combines PostgreSQL and Cassandra databases with Kafka queue service. It is recommended if you plan to manage 1M+ devices in production or handle high data ingestion rate (more than 5000 msg/sec). + +Create a docker compose file for the ThingsBoard Edge service: + +##### In Memory ```bash nano docker-compose.yml @@ -20,25 +37,22 @@ version: '3.8' services: mytbedge: restart: always - image: "thingsboard/tb-edge:${TB_EDGE_VERSION}" + image: "thingsboard/tb-edge:3.9.1EDGE" ports: - "8080:8080" - "1883:1883" - "5683-5688:5683-5688/udp" environment: SPRING_DATASOURCE_URL: jdbc:postgresql://postgres:5432/tb-edge - CLOUD_ROUTING_KEY: ${CLOUD_ROUTING_KEY} - CLOUD_ROUTING_SECRET: ${CLOUD_ROUTING_SECRET} - CLOUD_RPC_HOST: ${BASE_URL} - CLOUD_RPC_PORT: ${CLOUD_RPC_PORT} - CLOUD_RPC_SSL_ENABLED: ${CLOUD_RPC_SSL_ENABLED} + CLOUD_ROUTING_KEY: PUT_YOUR_EDGE_KEY_HERE # e.g. 19ea7ee8-5e6d-e642-4f32-05440a529015 + CLOUD_ROUTING_SECRET: PUT_YOUR_EDGE_SECRET_HERE # e.g. bztvkvfqsye7omv9uxlp + CLOUD_RPC_HOST: PUT_YOUR_CLOUD_IP # e.g. 192.168.1.1 or demo.thingsboard.io volumes: - tb-edge-data:/data - tb-edge-logs:/var/log/tb-edge - ${EXTRA_HOSTS} postgres: restart: always - image: "postgres:16" + image: "postgres:15" ports: - "5432" environment: @@ -58,24 +72,20 @@ volumes: ``` ##### [Optional] Update bind ports -If ThingsBoard Edge is going to be running on the same machine where ThingsBoard server (cloud) is running, you'll need to update docker compose port mapping to avoid port collision between ThingsBoard server and ThingsBoard Edge. +If ThingsBoard Edge is set to run on the same machine where the ThingsBoard server is operating, you need to update port configuration to prevent port collision between the ThingsBoard server and ThingsBoard Edge. -Please update next lines of `docker-compose.yml` file: +Ensure that the ports 18080, 11883, 15683-15688 are not used by any other application. -```text -ports: - - "18080:8080" - - "11883:1883" - - "15683-15688:5683-5688/udp" +Then, update the port configuration in the docker-compose.yml file: +```bash +sed -i ‘s/8080:8080/18080:8080/; s/1883:1883/11883:1883/; s/5683-5688:5683-5688\/udp/15683-15688:5683-5688\/udp/’ docker-compose.yml +{:copy-code} ``` -Make sure that ports above (18080, 11883, 15683-15688) are not used by any other application. - #### Start ThingsBoard Edge -Set the terminal in the directory which contains the `docker-compose.yml` file and execute the following commands to up this docker compose directly: +Set the terminal in the directory which contains the docker-compose.yml file and execute the following commands to up this docker compose directly: ```bash -docker compose up -d -docker compose logs -f mytbedge +docker compose up -d && docker compose logs -f mytbedge {:copy-code} ``` @@ -90,11 +100,12 @@ docker-compose up -d docker-compose logs -f mytbedge ``` -#### Open ThingsBoard Edge UI +#### Step 3. Open ThingsBoard Edge UI -Once started, you will be able to open **ThingsBoard Edge UI** using the following link http://localhost:8080. +Once the Edge service is started, open the Edge UI at http://localhost:8080. ###### NOTE: Edge HTTP bind port update -Use next **ThingsBoard Edge UI** link **http://localhost:18080** if you updated HTTP 8080 bind port to **18080**. +If the Edge HTTP bind port was changed to 18080 during Edge installation, access the ThingsBoard Edge instance at http://localhost:18080. +Please use your tenant credentials from local Server instance or ThingsBoard Live Demo to log in to the ThingsBoard Edge. diff --git a/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md b/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md index 992b5e2ee2..9d685c576a 100644 --- a/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md @@ -1,51 +1,48 @@ Here is the list of commands, that can be used to quickly install ThingsBoard Edge on Ubuntu Server and connect to the server. -#### Install Java 17 (OpenJDK) -ThingsBoard service is running on Java 17. Follow these instructions to install OpenJDK 17: +#### Step 1. Install Java 17 (OpenJDK) +ThingsBoard service is running on Java 17. To install OpenJDK 17, follow these instructions: ```bash -sudo apt update -sudo apt install openjdk-17-jdk +sudo apt update && sudo apt install openjdk-17-jdk {:copy-code} ``` -Please don't forget to configure your operating system to use OpenJDK 17 by default. -You can configure which version is the default using the following command: +Configure your operating system to use OpenJDK 17 by default. You can configure the default version by running the following command: ```bash sudo update-alternatives --config java {:copy-code} ``` -You can check the installation using the following command: +To check the installed Java version on your system, use the following command: ```bash java -version {:copy-code} ``` -Expected command output is: +The expected result is: ```text -openjdk version "17.x.xx" +openjdk version "17.x.xx" OpenJDK Runtime Environment (...) -OpenJDK 64-Bit Server VM (build ...) +OpenJDK 64-Bit Server VM (...) ``` -#### Configure PostgreSQL -ThingsBoard Edge uses PostgreSQL database as a local storage. -Instructions listed below will help you to install PostgreSQL. +#### Step 2. Configure ThingsBoard Edge Database -```bash -# install **wget** if not already installed: -sudo apt install -y wget +ThingsBoard Edge supports SQL and hybrid database approaches. See the architecture page for details. -# import the repository signing key: -wget --quiet -O - https://www.postgresql.org/media/keys/ACCC4CF8.asc | sudo apt-key add - +### Configure PostgreSQL +ThingsBoard Edge uses PostgreSQL database as a local storage. -# add repository contents to your system: -RELEASE=$(lsb_release -cs) -echo "deb http://apt.postgresql.org/pub/repos/apt/ ${RELEASE}"-pgdg main | sudo tee /etc/apt/sources.list.d/pgdg.list +To install the PostgreSQL database, run these commands: + +```bash +# Automated repository configuration: +sudo apt install -y postgresql-common +sudo /usr/share/postgresql-common/pgdg/apt.postgresql.org.sh # install and launch the postgresql service: sudo apt update @@ -54,25 +51,35 @@ sudo service postgresql start {:copy-code} ``` -Once PostgreSQL is installed you may want to create a new user or set the password for the main user. -The instructions below will help to set the password for main PostgreSQL user: +Once PostgreSQL is installed, it is recommended to set the password for the PostgreSQL main user. -```text -sudo su - postgres -psql -\password -\q +The following command will switch the current user to the PostgreSQL user and set the password directly in PostgreSQL. + +```bash +sudo -u postgres psql -c "\password" +{:copy-code} ``` -Then, press “Ctrl+D” to return to main user console and connect to the database to create ThingsBoard Edge DB: +Then, enter and confirm the password. -```text -psql -U postgres -d postgres -h 127.0.0.1 -W -CREATE DATABASE tb_edge; -\q +Finally, create a new PostgreSQL database named tb_edge by running the following command: + +```bash +echo "CREATE DATABASE tb_edge;" | psql -U postgres -d postgres -h 127.0.0.1 -W +{:copy-code} ``` -#### Thingsboard Edge service installation +#### Step 3. Choose Queue Service + +ThingsBoard Edge can use different messaging systems and brokers for storing messages and enabling communication between its services. Choose the appropriate queue implementation based on your specific business needs: + +In Memory: The built-in and default queue implementation. It is useful for development or proof-of-concept (PoC) environments, but is not recommended for production or any type of clustered deployments due to limited scalability. + +Kafka: Recommended for production deployments. This queue is used in the most of ThingsBoard production environments now. + +In Memory queue is built in and enabled by default. No additional configuration is required. + +#### Step 4. ThingsBoard Edge Service Installation Download installation package: ```bash @@ -87,39 +94,34 @@ sudo dpkg -i tb-edge-${TB_EDGE_TAG}.deb {:copy-code} ``` -#### Configure ThingsBoard Edge +#### Step 5. Configure ThingsBoard Edge To configure ThingsBoard Edge, you can use the following command to automatically update the configuration file with specific values: ```bash sudo sh -c 'cat <> /etc/tb-edge/conf/tb-edge.conf -export CLOUD_ROUTING_KEY=${CLOUD_ROUTING_KEY} -export CLOUD_ROUTING_SECRET=${CLOUD_ROUTING_SECRET} -export CLOUD_RPC_HOST=${BASE_URL} -export CLOUD_RPC_PORT=${CLOUD_RPC_PORT} -export CLOUD_RPC_SSL_ENABLED=${CLOUD_RPC_SSL_ENABLED} +export CLOUD_ROUTING_KEY= +export CLOUD_ROUTING_SECRET= +export CLOUD_RPC_HOST=demo.thingsboard.io +export CLOUD_RPC_PORT=7070 +export CLOUD_RPC_SSL_ENABLED=false EOL' {:copy-code} ``` -##### [Optional] Database Configuration -In case you changed default PostgreSQL datasource settings (**postgres**/**postgres**) please update the configuration file (**/etc/tb-edge/conf/tb-edge.conf**) with your actual values: - -```bash -sudo nano /etc/tb-edge/conf/tb-edge.conf -{:copy-code} -``` - -Please update the following lines in your configuration file. Make sure **to replace**: -- Replace 'postgres' with your actual PostgreSQL username; -- Replace 'PUT_YOUR_POSTGRESQL_PASSWORD_HERE' with your actual PostgreSQL password. +##### [Optional] Configure PostgreSQL +If you changed PostgreSQL default datasource settings, use the following command: ```bash +sudo sh -c 'cat <> /etc/tb-edge/conf/tb-edge.conf export SPRING_DATASOURCE_URL=jdbc:postgresql://localhost:5432/tb_edge export SPRING_DATASOURCE_USERNAME=postgres -export SPRING_DATASOURCE_PASSWORD=PUT_YOUR_POSTGRESQL_PASSWORD_HERE +export SPRING_DATASOURCE_PASSWORD= +EOL' {:copy-code} ``` +PUT_YOUR_POSTGRESQL_PASSWORD_HERE: Replace with your actual PostgreSQL user password. + ##### [Optional] Update bind ports If ThingsBoard Edge is going to be running on the same machine where ThingsBoard server (cloud) is running, you'll need to update configuration parameters to avoid port collision between ThingsBoard server and ThingsBoard Edge. @@ -138,7 +140,7 @@ EOL' Make sure that ports above (18080, 11883, 15683) are not used by any other application. -#### Run installation script +#### Step 6. Run installation Script Once ThingsBoard Edge is installed and configured please execute the following install script: @@ -147,14 +149,14 @@ sudo /usr/share/tb-edge/bin/install/install.sh {:copy-code} ``` -#### Restart ThingsBoard Edge service +#### Step 7. Restart ThingsBoard Edge Service ```bash sudo service tb-edge restart {:copy-code} ``` -#### Open ThingsBoard Edge UI +#### Step 8. Open ThingsBoard Edge UI Once started, you will be able to open **ThingsBoard Edge UI** using the following link http://localhost:8080. From dd493d1f52d4b1f254e947490bb6c297e689919b Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 25 Apr 2025 15:32:03 +0300 Subject: [PATCH 03/28] added tests --- .../cf/ctx/state/CalculatedFieldCtx.java | 4 +- .../ctx/state/SimpleCalculatedFieldState.java | 6 +- .../cf/CalculatedFieldIntegrationTest.java | 82 +++++++++++++++++++ .../SimpleCalculatedFieldConfiguration.java | 2 +- .../script/api/tbel/TbelCfCtx.java | 4 +- 5 files changed, 90 insertions(+), 8 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index 25083ed415..a3fdae319d 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -62,7 +62,7 @@ public class CalculatedFieldCtx { private final List argNames; private Output output; private String expression; - private boolean preserveLatestTs; + private boolean preserveMsgTs; private TbelInvokeService tbelInvokeService; private CalculatedFieldScriptEngine calculatedFieldScriptEngine; private ThreadLocal customExpression; @@ -96,7 +96,7 @@ public class CalculatedFieldCtx { this.argNames = new ArrayList<>(arguments.keySet()); this.output = configuration.getOutput(); this.expression = configuration.getExpression(); - this.preserveLatestTs = CalculatedFieldType.SIMPLE.equals(calculatedField.getType()) && ((SimpleCalculatedFieldConfiguration) configuration).isPreserveLastUpdateTs(); + this.preserveMsgTs = CalculatedFieldType.SIMPLE.equals(calculatedField.getType()) && ((SimpleCalculatedFieldConfiguration) configuration).isPreserveMsgTs(); this.tbelInvokeService = tbelInvokeService; this.maxDataPointsPerRollingArg = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 2ad961f72c..01091d2999 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -68,7 +68,7 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { Output output = ctx.getOutput(); Object result = formatResult(expressionResult, output.getDecimalsByDefault()); - JsonNode outputResult = createResultJson(ctx.isPreserveLatestTs(), output.getName(), result); + JsonNode outputResult = createResultJson(ctx.isPreserveMsgTs(), output.getName(), result); return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), outputResult)); } @@ -82,12 +82,12 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { : TbUtils.toFixed(expressionResult, decimals); } - private JsonNode createResultJson(boolean preserveLatestTs, String outputName, Object result) { + private JsonNode createResultJson(boolean preserveMsgTs, String outputName, Object result) { ObjectNode valuesNode = JacksonUtil.newObjectNode(); valuesNode.set(outputName, JacksonUtil.valueToTree(result)); long lastTimestamp = getLastUpdateTimestamp(); - if (preserveLatestTs && lastTimestamp != -1) { + if (preserveMsgTs && lastTimestamp != -1) { ObjectNode resultNode = JacksonUtil.newObjectNode(); resultNode.put("ts", lastTimestamp); resultNode.set("values", valuesNode); diff --git a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java index 002f668fdc..f65f6bc629 100644 --- a/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/CalculatedFieldIntegrationTest.java @@ -32,6 +32,7 @@ import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; +import org.thingsboard.server.common.data.cf.configuration.ScriptCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.debug.DebugSettings; import org.thingsboard.server.common.data.id.AssetProfileId; @@ -462,6 +463,87 @@ public class CalculatedFieldIntegrationTest extends CalculatedFieldControllerTes }); } + @Test + public void testSimpleCalculatedFieldWhenPreserveMsgTsIsTrue() throws Exception { + Device testDevice = createDevice("Test device", "1234567890"); + long ts = System.currentTimeMillis() - 300000L; + doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode(String.format("{\"ts\": %s, \"values\": {\"temperature\":30}}", ts))); + + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setEntityId(testDevice.getId()); + calculatedField.setType(CalculatedFieldType.SIMPLE); + calculatedField.setName("C to F"); + calculatedField.setDebugSettings(DebugSettings.all()); + calculatedField.setConfigurationVersion(1); + + SimpleCalculatedFieldConfiguration config = new SimpleCalculatedFieldConfiguration(); + + Argument argument = new Argument(); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); + config.setArguments(Map.of("T", argument)); + config.setExpression("(T * 9/5) + 32"); + + Output output = new Output(); + output.setName("fahrenheitTemp"); + output.setType(OutputType.TIME_SERIES); + config.setOutput(output); + + config.setPreserveMsgTs(true); + + calculatedField.setConfiguration(config); + + CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); + + await().alias("create CF -> perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + }); + } + + @Test + public void testScriptCalculatedFieldWhenUsedMsgTsInScript() throws Exception { + Device testDevice = createDevice("Test device", "1234567890"); + long ts = System.currentTimeMillis() - 300000L; + doPost("/api/plugins/telemetry/DEVICE/" + testDevice.getUuidId() + "/timeseries/" + DataConstants.SERVER_SCOPE, JacksonUtil.toJsonNode(String.format("{\"ts\": %s, \"values\": {\"temperature\":30}}", ts))); + + CalculatedField calculatedField = new CalculatedField(); + calculatedField.setEntityId(testDevice.getId()); + calculatedField.setType(CalculatedFieldType.SCRIPT); + calculatedField.setName("C to F"); + calculatedField.setDebugSettings(DebugSettings.all()); + calculatedField.setConfigurationVersion(1); + + ScriptCalculatedFieldConfiguration config = new ScriptCalculatedFieldConfiguration(); + + Argument argument = new Argument(); + ReferencedEntityKey refEntityKey = new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null); + argument.setRefEntityKey(refEntityKey); + config.setArguments(Map.of("T", argument)); + config.setExpression("return {\"ts\": ctx.msgTs, \"values\": {\"fahrenheitTemp\": (T * 1.8) + 32}};"); + + Output output = new Output(); + output.setType(OutputType.TIME_SERIES); + config.setOutput(output); + + calculatedField.setConfiguration(config); + + CalculatedField savedCalculatedField = doPost("/api/calculatedField", calculatedField, CalculatedField.class); + + await().alias("create CF -> perform initial calculation").atMost(TIMEOUT, TimeUnit.SECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode fahrenheitTemp = getLatestTelemetry(testDevice.getId(), "fahrenheitTemp"); + assertThat(fahrenheitTemp).isNotNull(); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("ts").asText()).isEqualTo(Long.toString(ts)); + assertThat(fahrenheitTemp.get("fahrenheitTemp").get(0).get("value").asText()).isEqualTo("86.0"); + }); + } + private ObjectNode getLatestTelemetry(EntityId entityId, String... keys) throws Exception { return doGetAsync("/api/plugins/telemetry/" + entityId.getEntityType() + "/" + entityId.getId() + "/values/timeseries?keys=" + String.join(",", keys), ObjectNode.class); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java index c748afdac4..af3cb4d5cd 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/SimpleCalculatedFieldConfiguration.java @@ -23,7 +23,7 @@ import org.thingsboard.server.common.data.cf.CalculatedFieldType; @EqualsAndHashCode(callSuper = true) public class SimpleCalculatedFieldConfiguration extends BaseCalculatedFieldConfiguration implements CalculatedFieldConfiguration { - private boolean preserveLastUpdateTs; + private boolean preserveMsgTs; @Override public CalculatedFieldType getType() { diff --git a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java index 1a8610cbab..7515cb5269 100644 --- a/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java +++ b/common/script/script-api/src/main/java/org/thingsboard/script/api/tbel/TbelCfCtx.java @@ -25,11 +25,11 @@ public class TbelCfCtx implements TbelCfObject { @Getter private final Map args; @Getter - private final long lastTs; + private final long msgTs; public TbelCfCtx(Map args, long lastUpdateTs) { this.args = Collections.unmodifiableMap(args); - this.lastTs = lastUpdateTs != -1 ? lastUpdateTs : System.currentTimeMillis(); + this.msgTs = lastUpdateTs != -1 ? lastUpdateTs : System.currentTimeMillis(); } @Override From 34592ef4b428d11dc7240354e5fabb20cb4af604 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 25 Apr 2025 16:49:29 +0300 Subject: [PATCH 04/28] added description --- .../ctx/state/BaseCalculatedFieldState.java | 6 +- .../en_US/calculated-field/expression_fn.md | 127 +++++++++++++----- 2 files changed, 96 insertions(+), 37 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index 84c61661ae..35879ac9aa 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -111,13 +111,11 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { private void updateLastUpdateTimestamp(ArgumentEntry entry) { if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - long ts = singleValueArgumentEntry.getTs(); - this.lastUpdateTimestamp = Math.max(this.lastUpdateTimestamp, ts); + this.lastUpdateTimestamp = singleValueArgumentEntry.getTs(); } else if (entry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { Map.Entry lastEntry = tsRollingArgumentEntry.getTsRecords().pollLastEntry(); if (lastEntry != null) { - long ts = lastEntry.getKey(); - this.lastUpdateTimestamp = Math.max(this.lastUpdateTimestamp, ts); + this.lastUpdateTimestamp = lastEntry.getKey(); } } } diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md b/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md index 373c50e2c6..9a988df042 100644 --- a/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md +++ b/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md @@ -1,7 +1,7 @@ ## Calculated Field TBEL Script Function The **calculate()** function is a user-defined script that enables custom calculations using [TBEL](${siteBaseUrl}/docs${docPlatformPrefix}/user-guide/tbel/) on telemetry and attribute data. -It receives arguments configured in the calculated field setup, along with an additional `ctx` object that provides access to all arguments. +It receives arguments configured in the calculated field setup, along with an additional `ctx` object that stores `msgTs` and provides access to all arguments. ### Function Signature @@ -44,7 +44,7 @@ Let's modify the function that converts Fahrenheit to Celsius to also return the var temperatureC = (temperatureF - 32) / 1.8; return { "ts": ctx.args.temperatureF.ts, - "values": { "temperatureC": toFixed(temperatureC, 2) } + "values": {"temperatureC": toFixed(temperatureC, 2)} }; ``` @@ -60,10 +60,22 @@ These contain time series data within a defined time window. Example format: "endTs": 1740644662896 }, "values": [ - { "ts": 1740644350000, "value": 72.32 }, - { "ts": 1740644360000, "value": 72.86 }, - { "ts": 1740644370000, "value": 73.58 }, - { "ts": 1740644380000, "value": "NaN" } + { + "ts": 1740644350000, + "value": 72.32 + }, + { + "ts": 1740644360000, + "value": 72.86 + }, + { + "ts": 1740644370000, + "value": 73.58 + }, + { + "ts": 1740644380000, + "value": "NaN" + } ] } } @@ -81,14 +93,18 @@ var firstItemTs = firstItem.ts; var firstItemValue = firstItem.value; var sum = 0.0; // iterate through all values and calculate the sum using foreach: -foreach(t: temperature) { - if(!isNaN(t.value)) { // check that the value is a valid number; +foreach(t +: +temperature +) +{ + if (!isNaN(t.value)) { // check that the value is a valid number; sum += t.value; } } // iterate through all values and calculate the sum using for loop: sum = 0.0; -for(var i = 0; i < temperature.values.size; i++) { +for (var i = 0; i < temperature.values.size; i++) { sum += temperature.values[i].value; } // use built-in function to calculate the sum @@ -146,12 +162,13 @@ function calculate(ctx, altitude, temperature) { Time series rolling arguments can be **merged** to align timestamps across multiple datasets. -| Method | Description | Returns | Example | -|:-----------------------------|:--------------------------------------------------------------------------------------------------------------------------|:----------------------------------------------------|:-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------| -| `merge(other, settings)` | Merges with another rolling argument. Aligns timestamps and filling missing values with the previous available value. | Merged object with `timeWindow` and aligned values. |

| -| `mergeAll(others, settings)` | Merges multiple rolling arguments. Aligns timestamps and filling missing values with the previous available value. | Merged object with `timeWindow` and aligned values. |

| +| Method | Description | Returns | Example | +|:-----------------------------|:----------------------------------------------------------------------------------------------------------------------|:----------------------------------------------------|:-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| `merge(other, settings)` | Merges with another rolling argument. Aligns timestamps and filling missing values with the previous available value. | Merged object with `timeWindow` and aligned values. |

| +| `mergeAll(others, settings)` | Merges multiple rolling arguments. Aligns timestamps and filling missing values with the previous available value. | Merged object with `timeWindow` and aligned values. |

| ##### Parameters + | Parameter | Description | |:---------------------|:-------------------------------------------------------------------------------------------------------------------------------------------------------------------------| | `other` or `others` | Another rolling argument or array of rolling arguments to merge with. | @@ -166,7 +183,11 @@ function calculate(ctx, temperature, defrost) { var merged = temperature.merge(defrost); var result = []; - foreach(item: merged) { + foreach(item +: + merged +) + { if (item.v1 > -5.0 && item.v2 == 0) { result.add({ ts: item.ts, @@ -187,29 +208,51 @@ function calculate(ctx, temperature, defrost) { The result is a list of issues that may be used to configure alarm rules: ```json -[{ +[ + { "ts": 1741613833843, "values": { - "issue": { - "temperature": -3.12, - "defrostState": false - } + "issue": { + "temperature": -3.12, + "defrostState": false + } } -}, { + }, + { "ts": 1741613923848, "values": { - "issue": { - "temperature": -4.16, - "defrostState": false - } + "issue": { + "temperature": -4.16, + "defrostState": false + } } -}] + } +] ``` ### Function return format The return format depends on the output type configured in the calculated field settings (default: **Time Series**). +### Message timestamp + +The `ctx` object also includes property `msgTs`, which represents the timestamp of the incoming telemetry message that triggered the calculated field execution in milliseconds. + +You can use `ctx.msgTs` to set the timestamp of the resulting output explicitly when returning a time series object. + +```javascript +var temperatureC = (temperatureF - 32) / 1.8; +return { + ts: ctx.msgTs, + values: { + "temperatureC": toFixed(temperatureC, 2) + } +} + +``` + +This ensures that the calculated data point aligns with the timestamp of the triggering telemetry. + ##### Time Series Output The function must return a JSON object or array with or without a timestamp. @@ -225,8 +268,14 @@ Without timestamp: "hvacState": "IDLE", "configuration": { "someNumber": 42, - "someArray": [1,2,3], - "someNestedObject": {"key": "value"} + "someArray": [ + 1, + 2, + 3 + ], + "someNestedObject": { + "key": "value" + } } } ``` @@ -243,10 +292,16 @@ With timestamp: "hvacState": "IDLE", "configuration": { "someNumber": 42, - "someArray": [1,2,3], - "someNestedObject": {"key": "value"} + "someArray": [ + 1, + 2, + 3 + ], + "someNestedObject": { + "key": "value" + } } - } + } } ``` @@ -265,7 +320,7 @@ Array containing multiple timestamps and different values of the `airDensity` : "values": { "airDensity": 1.07 } - } + } ] ``` @@ -282,8 +337,14 @@ Example below return 5 data points: airDensity (double), humidity (integer), hva "hvacState": "IDLE", "configuration": { "someNumber": 42, - "someArray": [1,2,3], - "someNestedObject": {"key": "value"} + "someArray": [ + 1, + 2, + 3 + ], + "someNestedObject": { + "key": "value" + } } } ``` From 644bc0c64799e68c427a797b6d49aa62d6874d38 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 28 Apr 2025 13:01:15 +0300 Subject: [PATCH 05/28] updated formatting result --- .../service/cf/ctx/state/SimpleCalculatedFieldState.java | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 01091d2999..d0eba5031c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -77,9 +77,10 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { if (decimals == null) { return expressionResult; } - return decimals.equals(0) - ? TbUtils.toInt(expressionResult) - : TbUtils.toFixed(expressionResult, decimals); + if (decimals.equals(0)) { + return TbUtils.toInt(expressionResult); + } + return TbUtils.toFixed(expressionResult, decimals); } private JsonNode createResultJson(boolean preserveMsgTs, String outputName, Object result) { From 1fff7f12715e6508b13162d54157d63da9fd042b Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 28 Apr 2025 14:48:36 +0300 Subject: [PATCH 06/28] updated script help page --- .../en_US/calculated-field/expression_fn.md | 64 ++++--------------- 1 file changed, 13 insertions(+), 51 deletions(-) diff --git a/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md b/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md index 9a988df042..4c54b01499 100644 --- a/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md +++ b/ui-ngx/src/assets/help/en_US/calculated-field/expression_fn.md @@ -60,22 +60,10 @@ These contain time series data within a defined time window. Example format: "endTs": 1740644662896 }, "values": [ - { - "ts": 1740644350000, - "value": 72.32 - }, - { - "ts": 1740644360000, - "value": 72.86 - }, - { - "ts": 1740644370000, - "value": 73.58 - }, - { - "ts": 1740644380000, - "value": "NaN" - } + { "ts": 1740644350000, "value": 72.32 }, + { "ts": 1740644360000, "value": 72.86 }, + { "ts": 1740644370000, "value": 73.58 }, + { "ts": 1740644380000, "value": "NaN" } ] } } @@ -93,12 +81,8 @@ var firstItemTs = firstItem.ts; var firstItemValue = firstItem.value; var sum = 0.0; // iterate through all values and calculate the sum using foreach: -foreach(t -: -temperature -) -{ - if (!isNaN(t.value)) { // check that the value is a valid number; +foreach(t: temperature) { + if(!isNaN(t.value)) { // check that the value is a valid number; sum += t.value; } } @@ -183,11 +167,7 @@ function calculate(ctx, temperature, defrost) { var merged = temperature.merge(defrost); var result = []; - foreach(item -: - merged -) - { + foreach(item: merged) { if (item.v1 > -5.0 && item.v2 == 0) { result.add({ ts: item.ts, @@ -268,14 +248,8 @@ Without timestamp: "hvacState": "IDLE", "configuration": { "someNumber": 42, - "someArray": [ - 1, - 2, - 3 - ], - "someNestedObject": { - "key": "value" - } + "someArray": [1,2,3], + "someNestedObject": {"key": "value"} } } ``` @@ -292,14 +266,8 @@ With timestamp: "hvacState": "IDLE", "configuration": { "someNumber": 42, - "someArray": [ - 1, - 2, - 3 - ], - "someNestedObject": { - "key": "value" - } + "someArray": [1,2,3], + "someNestedObject": {"key": "value"} } } } @@ -337,14 +305,8 @@ Example below return 5 data points: airDensity (double), humidity (integer), hva "hvacState": "IDLE", "configuration": { "someNumber": 42, - "someArray": [ - 1, - 2, - 3 - ], - "someNestedObject": { - "key": "value" - } + "someArray": [1,2,3], + "someNestedObject": {"key": "value"} } } ``` From 08679cf561c24842735b61efae35f85253e535a7 Mon Sep 17 00:00:00 2001 From: yevhenii Date: Tue, 29 Apr 2025 11:33:50 +0300 Subject: [PATCH 07/28] PostgreSQL View for active Edges - Moved filtering logic for active edges from inline native query to a database view (edge_active_attribute_view) to simplify the query and improve maintainability. --- .../server/dao/sql/edge/EdgeRepository.java | 13 ++++--------- .../sql/schema-views-and-functions.sql | 19 +++++++++++++++++++ .../resources/sql/psql/drop-all-tables.sql | 1 + 3 files changed, 24 insertions(+), 9 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java index bd1dda54b8..66d65cf602 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java @@ -44,15 +44,10 @@ public interface EdgeRepository extends JpaRepository { "WHERE d.id = :edgeId") EdgeInfoEntity findEdgeInfoById(@Param("edgeId") UUID edgeId); - @Query(value = "SELECT ee.id, ee.created_time, ee.additional_info, ee.customer_id, " + - "ee.root_rule_chain_id, ee.type, ee.name, ee.label, ee.routing_key, " + - "ee.secret, ee.tenant_id, ee.version " + - "FROM edge ee " + - "JOIN attribute_kv ON ee.id = attribute_kv.entity_id " + - "JOIN key_dictionary ON attribute_kv.attribute_key = key_dictionary.key_id " + - "WHERE attribute_kv.bool_v = true AND key_dictionary.key = 'active' " + - "AND (:textSearch IS NULL OR ee.name ILIKE CONCAT('%', :textSearch, '%')) " + - "ORDER BY ee.id", nativeQuery = true) + @Query(value = "SELECT * " + + "FROM edge_active_attribute_view edge_active " + + "WHERE (:textSearch IS NULL OR edge_active.name ILIKE CONCAT('%', :textSearch, '%')) " + + "ORDER BY edge_active.id", nativeQuery = true) Page findActiveEdges(@Param("textSearch") String textSearch, Pageable pageable); diff --git a/dao/src/main/resources/sql/schema-views-and-functions.sql b/dao/src/main/resources/sql/schema-views-and-functions.sql index a1abef81ca..b65dc41a90 100644 --- a/dao/src/main/resources/sql/schema-views-and-functions.sql +++ b/dao/src/main/resources/sql/schema-views-and-functions.sql @@ -72,6 +72,25 @@ u.first_name as assignee_first_name, u.last_name as assignee_last_name, u.email FROM alarm a LEFT JOIN tb_user u ON u.id = a.assignee_id; +DROP VIEW IF EXISTS edge_active_attribute_view CASCADE; +CREATE OR REPLACE VIEW edge_active_attribute_view AS +SELECT ee.id + , ee.created_time + , ee.additional_info + , ee.customer_id + , ee.root_rule_chain_id + , ee.type + , ee.name + , ee.label + , ee.routing_key + , ee.secret + , ee.tenant_id + , ee.version +FROM edge ee + JOIN attribute_kv ON ee.id = attribute_kv.entity_id + JOIN key_dictionary ON attribute_kv.attribute_key = key_dictionary.key_id +WHERE attribute_kv.bool_v = true AND key_dictionary.key = 'active'; + CREATE OR REPLACE FUNCTION create_or_update_active_alarm( t_id uuid, c_id uuid, a_id uuid, a_created_ts bigint, a_o_id uuid, a_o_type integer, a_type varchar, diff --git a/dao/src/test/resources/sql/psql/drop-all-tables.sql b/dao/src/test/resources/sql/psql/drop-all-tables.sql index da6eca161b..0387658f2d 100644 --- a/dao/src/test/resources/sql/psql/drop-all-tables.sql +++ b/dao/src/test/resources/sql/psql/drop-all-tables.sql @@ -14,6 +14,7 @@ DROP VIEW IF EXISTS device_info_active_attribute_view CASCADE; DROP VIEW IF EXISTS device_info_active_ts_view CASCADE; DROP VIEW IF EXISTS device_info_view CASCADE; DROP VIEW IF EXISTS alarm_info CASCADE; +DROP VIEW IF EXISTS edge_acitve_attribute_view CASCADE; DROP TABLE IF EXISTS admin_settings; DROP TABLE IF EXISTS entity_alarm; From 18916bf693164500f6bdbe317dd4167e2482bc4b Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 5 May 2025 08:32:31 +0300 Subject: [PATCH 08/28] added check for timestamp --- .../server/controller/CalculatedFieldController.java | 2 +- .../service/cf/ctx/state/BaseCalculatedFieldState.java | 6 ++---- 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java index 7e4e041600..a28146ab6e 100644 --- a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java +++ b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java @@ -279,7 +279,7 @@ public class CalculatedFieldController extends BaseController { lastUpdateTimestamp = Math.max(lastUpdateTimestamp, maxTs); } } - return lastUpdateTimestamp; + return lastUpdateTimestamp == -1 ? System.currentTimeMillis() : lastUpdateTimestamp; } private & HasTenantId, I extends EntityId> void checkReferencedEntities(CalculatedFieldConfiguration calculatedFieldConfig, SecurityUser user) throws ThingsboardException { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index 35879ac9aa..e4b03b4cab 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -113,10 +113,8 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { this.lastUpdateTimestamp = singleValueArgumentEntry.getTs(); } else if (entry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { - Map.Entry lastEntry = tsRollingArgumentEntry.getTsRecords().pollLastEntry(); - if (lastEntry != null) { - this.lastUpdateTimestamp = lastEntry.getKey(); - } + Map.Entry lastEntry = tsRollingArgumentEntry.getTsRecords().lastEntry(); + this.lastUpdateTimestamp = (lastEntry != null) ? lastEntry.getKey() : System.currentTimeMillis(); } } From 8dcd8f707d9a1686caa6c6b47754113a0026d5c3 Mon Sep 17 00:00:00 2001 From: Yevhenii Date: Tue, 6 May 2025 09:37:21 +0300 Subject: [PATCH 09/28] Update edge install instructions - changed instructions --- .../edge/instructions/install/docker/instructions.md | 10 ++++++---- .../edge/instructions/install/ubuntu/instructions.md | 10 +++++----- 2 files changed, 11 insertions(+), 9 deletions(-) diff --git a/application/src/main/data/json/edge/instructions/install/docker/instructions.md b/application/src/main/data/json/edge/instructions/install/docker/instructions.md index 4457962d86..3d42cffa8a 100644 --- a/application/src/main/data/json/edge/instructions/install/docker/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/docker/instructions.md @@ -37,16 +37,18 @@ version: '3.8' services: mytbedge: restart: always - image: "thingsboard/tb-edge:3.9.1EDGE" + image: "thingsboard/tb-edge:${TB_EDGE_VERSION}" ports: - "8080:8080" - "1883:1883" - "5683-5688:5683-5688/udp" environment: SPRING_DATASOURCE_URL: jdbc:postgresql://postgres:5432/tb-edge - CLOUD_ROUTING_KEY: PUT_YOUR_EDGE_KEY_HERE # e.g. 19ea7ee8-5e6d-e642-4f32-05440a529015 - CLOUD_ROUTING_SECRET: PUT_YOUR_EDGE_SECRET_HERE # e.g. bztvkvfqsye7omv9uxlp - CLOUD_RPC_HOST: PUT_YOUR_CLOUD_IP # e.g. 192.168.1.1 or demo.thingsboard.io + CLOUD_ROUTING_KEY: ${CLOUD_ROUTING_KEY} + CLOUD_ROUTING_SECRET: ${CLOUD_ROUTING_SECRET} + CLOUD_RPC_HOST: ${BASE_URL} + CLOUD_RPC_PORT: ${CLOUD_RPC_PORT} + CLOUD_RPC_SSL_ENABLED: ${CLOUD_RPC_SSL_ENABLED} volumes: - tb-edge-data:/data - tb-edge-logs:/var/log/tb-edge diff --git a/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md b/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md index 9d685c576a..6ac45d325a 100644 --- a/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md @@ -99,11 +99,11 @@ To configure ThingsBoard Edge, you can use the following command to automatical ```bash sudo sh -c 'cat <> /etc/tb-edge/conf/tb-edge.conf -export CLOUD_ROUTING_KEY= -export CLOUD_ROUTING_SECRET= -export CLOUD_RPC_HOST=demo.thingsboard.io -export CLOUD_RPC_PORT=7070 -export CLOUD_RPC_SSL_ENABLED=false +export CLOUD_ROUTING_KEY=${CLOUD_ROUTING_KEY} +export CLOUD_ROUTING_SECRET=${CLOUD_ROUTING_SECRET} +export CLOUD_RPC_HOST=${BASE_URL} +export CLOUD_RPC_PORT=${CLOUD_RPC_PORT} +export CLOUD_RPC_SSL_ENABLED=${CLOUD_RPC_SSL_ENABLED} EOL' {:copy-code} ``` From 81e35e134e8377f2e24f106aba1188a235aee153 Mon Sep 17 00:00:00 2001 From: Yevhenii Date: Wed, 7 May 2025 17:40:59 +0300 Subject: [PATCH 10/28] Validation DefaultEdgeRuleChain in DeviceProfile - added validation for DefaultEdgeRuleChain - refactoring --- .../validator/DeviceProfileDataValidator.java | 28 +++++++++++++------ 1 file changed, 19 insertions(+), 9 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java index a9e14229f8..e32195022c 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java @@ -43,6 +43,7 @@ import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.Abstrac import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.LwM2MBootstrapServerCredential; import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.RPKLwM2MBootstrapServerCredential; import org.thingsboard.server.common.data.device.profile.lwm2m.bootstrap.X509LwM2MBootstrapServerCredential; +import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.queue.Queue; import org.thingsboard.server.common.data.rule.RuleChain; @@ -188,13 +189,11 @@ public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator 65535){ + if (serverConfig.isBootstrapServerIs()) { + if (serverConfig.getShortServerId() < 0 || serverConfig.getShortServerId() > 65535) { throw new DeviceCredentialsValidationException("Bootstrap Server ShortServerId must be in range [0 - 65535]!"); } } else { @@ -423,4 +432,5 @@ public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator Date: Tue, 13 May 2025 15:19:24 +0300 Subject: [PATCH 11/28] Add queue prefix in msa tests --- .../test/java/org/thingsboard/server/msa/ContainerTestSuite.java | 1 + 1 file changed, 1 insertion(+) diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java index e00b0828c9..0573aee8fe 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java @@ -130,6 +130,7 @@ public class ContainerTestSuite { Map queueEnv = new HashMap<>(); queueEnv.put("TB_QUEUE_TYPE", QUEUE_TYPE); + queueEnv.put("TB_QUEUE_PREFIX", "test"); switch (QUEUE_TYPE) { case "kafka": composeFiles.add(new File(targetDir + "docker-compose.kafka.yml")); From ef65dd90263f982ff032ca835b91bd8d4e48e6a7 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 13 May 2025 16:41:57 +0300 Subject: [PATCH 12/28] Fix missing queue prefixes --- .../server/service/edqs/KafkaEdqsSyncService.java | 5 +++-- .../thingsboard/server/edqs/processor/EdqsProcessor.java | 8 +++++--- .../server/edqs/state/KafkaEdqsStateService.java | 6 ++++-- .../server/queue/discovery/HashPartitionService.java | 2 +- .../org/thingsboard/server/queue/kafka/TbKafkaAdmin.java | 3 +++ 5 files changed, 16 insertions(+), 8 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edqs/KafkaEdqsSyncService.java b/application/src/main/java/org/thingsboard/server/service/edqs/KafkaEdqsSyncService.java index 239fd9dc42..43b0c575a0 100644 --- a/application/src/main/java/org/thingsboard/server/service/edqs/KafkaEdqsSyncService.java +++ b/application/src/main/java/org/thingsboard/server/service/edqs/KafkaEdqsSyncService.java @@ -18,6 +18,7 @@ package org.thingsboard.server.service.edqs; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.edqs.EdqsConfig; import org.thingsboard.server.queue.kafka.TbKafkaAdmin; import org.thingsboard.server.queue.kafka.TbKafkaSettings; @@ -32,11 +33,11 @@ public class KafkaEdqsSyncService extends EdqsSyncService { private final boolean syncNeeded; - public KafkaEdqsSyncService(TbKafkaSettings kafkaSettings, EdqsConfig edqsConfig) { + public KafkaEdqsSyncService(TbKafkaSettings kafkaSettings, TopicService topicService, EdqsConfig edqsConfig) { TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, Collections.emptyMap()); this.syncNeeded = kafkaAdmin.areAllTopicsEmpty(IntStream.range(0, edqsConfig.getPartitions()) .mapToObj(partition -> TopicPartitionInfo.builder() - .topic(edqsConfig.getEventsTopic()) + .topic(topicService.buildTopicName(edqsConfig.getEventsTopic())) .partition(partition) .build().getFullTopicName()) .collect(Collectors.toSet())); diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java index 78d17ef368..72919c3fc0 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/processor/EdqsProcessor.java @@ -58,6 +58,7 @@ import org.thingsboard.server.queue.TbQueueResponseTemplate; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.common.consumer.PartitionedQueueConsumerManager; import org.thingsboard.server.queue.discovery.QueueKey; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.edqs.EdqsComponent; import org.thingsboard.server.queue.edqs.EdqsConfig; @@ -88,6 +89,7 @@ public class EdqsProcessor implements TbQueueHandler, private final EdqsRepository repository; private final EdqsConfig config; private final EdqsPartitionService partitionService; + private final TopicService topicService; private final ConfigurableApplicationContext applicationContext; private final EdqsStateService stateService; @@ -123,7 +125,7 @@ public class EdqsProcessor implements TbQueueHandler, eventConsumer = PartitionedQueueConsumerManager.>create() .queueKey(new QueueKey(ServiceType.EDQS, config.getEventsTopic())) - .topic(config.getEventsTopic()) + .topic(topicService.buildTopicName(config.getEventsTopic())) .pollInterval(config.getPollInterval()) .msgPackProcessor((msgs, consumer, config) -> { for (TbProtoQueueMsg queueMsg : msgs) { @@ -164,9 +166,9 @@ public class EdqsProcessor implements TbQueueHandler, try { Set newPartitions = event.getNewPartitions().get(new QueueKey(ServiceType.EDQS)); - stateService.process(withTopic(newPartitions, config.getStateTopic())); + stateService.process(withTopic(newPartitions, topicService.buildTopicName(config.getStateTopic()))); // eventsConsumer's partitions are updated by stateService - responseTemplate.subscribe(withTopic(newPartitions, config.getRequestsTopic())); // TODO: we subscribe to partitions before we are ready. implement consumer-per-partition version for request template + responseTemplate.subscribe(withTopic(newPartitions, topicService.buildTopicName(config.getRequestsTopic()))); // TODO: we subscribe to partitions before we are ready. implement consumer-per-partition version for request template Set oldPartitions = event.getOldPartitions().get(new QueueKey(ServiceType.EDQS)); if (CollectionsUtil.isNotEmpty(oldPartitions)) { diff --git a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java index 0efe6e7d3b..ddbdc3253a 100644 --- a/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java +++ b/common/edqs/src/main/java/org/thingsboard/server/edqs/state/KafkaEdqsStateService.java @@ -36,6 +36,7 @@ import org.thingsboard.server.queue.common.consumer.QueueConsumerManager; import org.thingsboard.server.queue.common.state.KafkaQueueStateService; import org.thingsboard.server.queue.common.state.QueueStateService; import org.thingsboard.server.queue.discovery.QueueKey; +import org.thingsboard.server.queue.discovery.TopicService; import org.thingsboard.server.queue.edqs.EdqsConfig; import org.thingsboard.server.queue.edqs.KafkaEdqsComponent; import org.thingsboard.server.queue.edqs.KafkaEdqsQueueFactory; @@ -59,6 +60,7 @@ public class KafkaEdqsStateService implements EdqsStateService { private final EdqsConfig config; private final EdqsPartitionService partitionService; private final KafkaEdqsQueueFactory queueFactory; + private final TopicService topicService; @Autowired @Lazy private EdqsProcessor edqsProcessor; @@ -78,7 +80,7 @@ public class KafkaEdqsStateService implements EdqsStateService { TbKafkaAdmin queueAdmin = queueFactory.getEdqsQueueAdmin(); stateConsumer = PartitionedQueueConsumerManager.>create() .queueKey(new QueueKey(ServiceType.EDQS, config.getStateTopic())) - .topic(config.getStateTopic()) + .topic(topicService.buildTopicName(config.getStateTopic())) .pollInterval(config.getPollInterval()) .msgPackProcessor((msgs, consumer, config) -> { for (TbProtoQueueMsg queueMsg : msgs) { @@ -176,7 +178,7 @@ public class KafkaEdqsStateService implements EdqsStateService { if (queueStateService.getPartitions().isEmpty()) { Set allPartitions = IntStream.range(0, config.getPartitions()) .mapToObj(partition -> TopicPartitionInfo.builder() - .topic(config.getEventsTopic()) + .topic(topicService.buildTopicName(config.getEventsTopic())) .partition(partition) .build()) .collect(Collectors.toSet()); diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java index 7186bf7055..ec5c675e32 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/HashPartitionService.java @@ -156,7 +156,7 @@ public class HashPartitionService implements PartitionService { @Override public String getTopic(QueueKey queueKey) { - return partitionTopicsMap.get(queueKey); + return topicService.buildTopicName(partitionTopicsMap.get(queueKey)); } private void doInitRuleEnginePartitions() { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java index 1e0064a5c8..6addce2865 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java @@ -136,6 +136,9 @@ public class TbKafkaAdmin implements TbQueueAdmin, TbEdgeQueueAdmin { } public CreateTopicsResult createTopic(NewTopic topic) { + if (!topic.name().startsWith("test.")) { // FIXME: remove me + log.error("Creating topic without configured prefix: {}", topic.name(), new RuntimeException("stacktrace")); + } return settings.getAdminClient().createTopics(Collections.singletonList(topic)); } From 95d1b5ebd7de3f4eaa63ae8d2e9262912d5c05d9 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 15 May 2025 11:17:31 +0300 Subject: [PATCH 13/28] Minor refactoring for TopicPartitionInfo --- .../server/service/queue/DefaultTbCoreConsumerService.java | 2 +- .../server/common/msg/queue/TopicPartitionInfo.java | 4 ---- .../java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java | 3 --- .../queue/usagestats/DefaultTbApiUsageReportClient.java | 2 +- 4 files changed, 2 insertions(+), 9 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 47a0c473e9..d3743eb084 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -251,7 +251,7 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService tpi.newByTopic(usageStatsConsumer.getConsumer().getTopic())) + .map(tpi -> tpi.withTopic(usageStatsConsumer.getConsumer().getTopic())) .collect(Collectors.toSet())); } diff --git a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java index b18debaf49..80eaede7bf 100644 --- a/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java +++ b/common/message/src/main/java/org/thingsboard/server/common/msg/queue/TopicPartitionInfo.java @@ -57,10 +57,6 @@ public class TopicPartitionInfo { this(topic, tenantId, partition, false, myPartition); } - public TopicPartitionInfo newByTopic(String topic) { - return new TopicPartitionInfo(topic, this.tenantId, this.partition, this.useInternalPartition, this.myPartition); - } - public String getTopic() { return topic; } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java index 6addce2865..1e0064a5c8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java @@ -136,9 +136,6 @@ public class TbKafkaAdmin implements TbQueueAdmin, TbEdgeQueueAdmin { } public CreateTopicsResult createTopic(NewTopic topic) { - if (!topic.name().startsWith("test.")) { // FIXME: remove me - log.error("Creating topic without configured prefix: {}", topic.name(), new RuntimeException("stacktrace")); - } return settings.getAdminClient().createTopics(Collections.singletonList(topic)); } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java b/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java index 715020dc7c..34543a07e8 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/usagestats/DefaultTbApiUsageReportClient.java @@ -131,7 +131,7 @@ public class DefaultTbApiUsageReportClient implements TbApiUsageReportClient { report.forEach((parent, statsMsg) -> { try { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, parent.getTenantId(), parent.getId()) - .newByTopic(msgProducer.getDefaultTopic()); + .withTopic(msgProducer.getDefaultTopic()); reportStatsPerTpi.computeIfAbsent(tpi, k -> new ArrayList<>()).add(statsMsg.build()); } catch (TenantNotFoundException e) { log.debug("Couldn't report usage stats for non-existing tenant: {}", e.getTenantId()); From 52f7683b322f3bfec1288343b9345b7d2a1785e1 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 15 May 2025 12:02:59 +0300 Subject: [PATCH 14/28] Cleanup ContainerTestSuite --- .../server/msa/ContainerTestSuite.java | 41 ++----------------- 1 file changed, 3 insertions(+), 38 deletions(-) diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java index 0573aee8fe..5a53862164 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java @@ -29,7 +29,6 @@ import java.nio.file.Path; import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; -import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.UUID; @@ -46,7 +45,6 @@ public class ContainerTestSuite { final static boolean IS_REDIS_SENTINEL = Boolean.parseBoolean(System.getProperty("blackBoxTests.redisSentinel")); final static boolean IS_REDIS_SSL = Boolean.parseBoolean(System.getProperty("blackBoxTests.redisSsl")); final static boolean IS_HYBRID_MODE = Boolean.parseBoolean(System.getProperty("blackBoxTests.hybridMode")); - final static String QUEUE_TYPE = System.getProperty("blackBoxTests.queue", "kafka"); private static final String SOURCE_DIR = "./../../docker/"; private static final String TB_CORE_LOG_REGEXP = ".*Starting polling for events.*"; private static final String TRANSPORTS_LOG_REGEXP = ".*Going to recalculate partitions.*"; @@ -122,45 +120,12 @@ public class ContainerTestSuite { new File(targetDir + (IS_HYBRID_MODE ? "docker-compose.hybrid.yml" : "docker-compose.postgres.yml")), new File(targetDir + (IS_HYBRID_MODE ? "docker-compose.hybrid-test-extras.yml" : "docker-compose.postgres-test-extras.yml")), new File(targetDir + "docker-compose.postgres.volumes.yml"), - new File(targetDir + "docker-compose." + QUEUE_TYPE + ".yml"), + new File(targetDir + "docker-compose.kafka.yml"), new File(targetDir + resolveRedisComposeFile()), new File(targetDir + resolveRedisComposeVolumesFile()), new File(targetDir + ("docker-selenium.yml")) )); - - Map queueEnv = new HashMap<>(); - queueEnv.put("TB_QUEUE_TYPE", QUEUE_TYPE); - queueEnv.put("TB_QUEUE_PREFIX", "test"); - switch (QUEUE_TYPE) { - case "kafka": - composeFiles.add(new File(targetDir + "docker-compose.kafka.yml")); - break; - case "aws-sqs": - replaceInFile(targetDir, "queue-aws-sqs.env", - Map.of("YOUR_KEY", getSysProp("blackBoxTests.awsKey"), - "YOUR_SECRET", getSysProp("blackBoxTests.awsSecret"), - "YOUR_REGION", getSysProp("blackBoxTests.awsRegion"))); - break; - case "rabbitmq": - composeFiles.add(new File(targetDir + "docker-compose.rabbitmq-server.yml")); - replaceInFile(targetDir, "queue-rabbitmq.env", - Map.of("localhost", "rabbitmq")); - break; - case "service-bus": - replaceInFile(targetDir, "queue-service-bus.env", - Map.of("YOUR_NAMESPACE_NAME", getSysProp("blackBoxTests.serviceBusNamespace"), - "YOUR_SAS_KEY_NAME", getSysProp("blackBoxTests.serviceBusSASPolicy"))); - replaceInFile(targetDir, "queue-service-bus.env", - Map.of("YOUR_SAS_KEY", getSysProp("blackBoxTests.serviceBusPrimaryKey"))); - break; - case "pubsub": - replaceInFile(targetDir, "queue-pubsub.env", - Map.of("YOUR_PROJECT_ID", getSysProp("blackBoxTests.pubSubProjectId"), - "YOUR_SERVICE_ACCOUNT", getSysProp("blackBoxTests.pubSubServiceAccount"))); - break; - default: - throw new RuntimeException("Unsupported queue type: " + QUEUE_TYPE); - } + addToFile(targetDir, "queue-kafka.env", Map.of("TB_QUEUE_PREFIX", "test")); if (IS_HYBRID_MODE) { composeFiles.add(new File(targetDir + "docker-compose.cassandra.volumes.yml")); @@ -172,7 +137,7 @@ public class ContainerTestSuite { .withOptions("--compatibility") .withTailChildContainers(!skipTailChildContainers) .withEnv(installTb.getEnv()) - .withEnv(queueEnv) + .withEnv("TB_QUEUE_TYPE", "kafka") .withEnv("LOAD_BALANCER_NAME", "") .withExposedService("haproxy", 80, Wait.forHttp("/swagger-ui.html").withStartupTimeout(CONTAINER_STARTUP_TIMEOUT)) .withExposedService("broker", 1883) From 7a0c2b7763e7645dab4b34e9c0c28c6ad08740c7 Mon Sep 17 00:00:00 2001 From: Vladyslav_Prykhodko Date: Thu, 15 May 2025 17:12:17 +0300 Subject: [PATCH 15/28] UI: Add to CF - Use message timestamp --- .../calculated-field-dialog.component.html | 13 ++++++++--- .../calculated-field-dialog.component.scss | 2 +- .../calculated-field-dialog.component.ts | 22 ++++++++++++++++++- .../shared/models/calculated-field.models.ts | 10 +++++++++ .../assets/locale/locale.constant-en_US.json | 4 +++- 5 files changed, 45 insertions(+), 6 deletions(-) diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.html b/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.html index 8e9baa6f90..60f33ff005 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.html +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.html @@ -160,8 +160,8 @@ } @if (fieldFormGroup.get('type').value === CalculatedFieldType.SIMPLE) { -
- +
+ {{ (outputFormGroup.get('type').value === OutputType.Timeseries ? 'calculated-fields.timeseries-key' @@ -181,7 +181,7 @@ } - + {{ 'calculated-fields.decimals-by-default' | translate }} @if (outputFormGroup.get('decimalsByDefault').errors && outputFormGroup.get('decimalsByDefault').touched) { @@ -189,6 +189,13 @@ }
+
+ +
+ calculated-fields.use-message-timestamp +
+
+
}
diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.scss b/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.scss index 8bc422eed1..efcd62efd4 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.scss +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.scss @@ -45,7 +45,7 @@ &-key { color: #c24c1a; } - &-time-window, &-values, &-func, &-value, &-ts { + &-time-window, &-values, &-func, &-value, &-ts, &-msgTs { color: #7214D0; } &-start-ts, &-end-ts { diff --git a/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.ts b/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.ts index 52051aa6f7..4aa4eca425 100644 --- a/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.ts +++ b/ui-ngx/src/app/modules/home/components/calculated-fields/components/dialog/calculated-field-dialog.component.ts @@ -77,6 +77,7 @@ export class CalculatedFieldDialogComponent extends DialogComponent Date: Mon, 19 May 2025 10:49:16 +0300 Subject: [PATCH 16/28] Add missing test queue prefix for EDQS --- .../test/java/org/thingsboard/server/msa/ContainerTestSuite.java | 1 + 1 file changed, 1 insertion(+) diff --git a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java index 5a53862164..345d0b7144 100644 --- a/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java +++ b/msa/black-box-tests/src/test/java/org/thingsboard/server/msa/ContainerTestSuite.java @@ -126,6 +126,7 @@ public class ContainerTestSuite { new File(targetDir + ("docker-selenium.yml")) )); addToFile(targetDir, "queue-kafka.env", Map.of("TB_QUEUE_PREFIX", "test")); + addToFile(targetDir, "tb-edqs.env", Map.of("TB_QUEUE_PREFIX", "test")); if (IS_HYBRID_MODE) { composeFiles.add(new File(targetDir + "docker-compose.cassandra.volumes.yml")); From e12b4015b1f568d7b3de12072fd52a9bee19a0a3 Mon Sep 17 00:00:00 2001 From: Artem Dzhereleiko Date: Tue, 20 May 2025 12:06:30 +0300 Subject: [PATCH 17/28] UI: Fixed HP curcuit breaker widget type fqn --- .../widget_bundles/high_performance_scada_energy_system.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/application/src/main/data/json/system/widget_bundles/high_performance_scada_energy_system.json b/application/src/main/data/json/system/widget_bundles/high_performance_scada_energy_system.json index a06ccc0819..bc14e14023 100644 --- a/application/src/main/data/json/system/widget_bundles/high_performance_scada_energy_system.json +++ b/application/src/main/data/json/system/widget_bundles/high_performance_scada_energy_system.json @@ -15,7 +15,7 @@ "hp_wind_turbine_cluster", "hp_fuel_generator", "hp_industrial_fuel_generator", - "hp_circuit_breaker2", + "hp_circuit_breaker", "hp_horizontal_circuit_breaker", "hp_voltage_relay", "hp_3_phase_voltage_relay", From 440087384d3d478011a33756952d0d061c20af2f Mon Sep 17 00:00:00 2001 From: dshvaika Date: Tue, 20 May 2025 12:27:26 +0300 Subject: [PATCH 18/28] deduplication node: fixed retry mechanism --- .../deduplication/TbMsgDeduplicationNode.java | 17 +++-- .../transform/TbMsgDeduplicationNodeTest.java | 76 +++++++++++++++++++ 2 files changed, 85 insertions(+), 8 deletions(-) diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java index f7a7c6e4dd..ce5fba102a 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/deduplication/TbMsgDeduplicationNode.java @@ -64,7 +64,7 @@ import static org.thingsboard.server.common.data.DataConstants.QUEUE_NAME; @Slf4j public class TbMsgDeduplicationNode implements TbNode { - public static final int TB_MSG_DEDUPLICATION_RETRY_DELAY = 10; + public static final long TB_MSG_DEDUPLICATION_RETRY_DELAY = 10L; private TbMsgDeduplicationNodeConfiguration config; @@ -217,16 +217,17 @@ public class TbMsgDeduplicationNode implements TbNode { } private void enqueueForTellNextWithRetry(TbContext ctx, TbMsg msg, int retryAttempt) { - if (config.getMaxRetries() > retryAttempt) { + if (retryAttempt <= config.getMaxRetries()) { ctx.enqueueForTellNext(msg, TbNodeConnectionType.SUCCESS, - () -> { - log.trace("[{}][{}][{}] Successfully enqueue deduplication result message!", ctx.getSelfId(), msg.getOriginator(), retryAttempt); - }, + () -> log.trace("[{}][{}][{}] Successfully enqueue deduplication result message!", ctx.getSelfId(), msg.getOriginator(), retryAttempt), throwable -> { log.trace("[{}][{}][{}] Failed to enqueue deduplication output message due to: ", ctx.getSelfId(), msg.getOriginator(), retryAttempt, throwable); - ctx.schedule(() -> { - enqueueForTellNextWithRetry(ctx, msg, retryAttempt + 1); - }, TB_MSG_DEDUPLICATION_RETRY_DELAY, TimeUnit.SECONDS); + if (retryAttempt < config.getMaxRetries()) { + ctx.schedule(() -> enqueueForTellNextWithRetry(ctx, msg, retryAttempt + 1), TB_MSG_DEDUPLICATION_RETRY_DELAY, TimeUnit.SECONDS); + } else { + log.trace("[{}][{}] Max retries [{}] exhausted. Dropping deduplication result message [{}]", + ctx.getSelfId(), msg.getOriginator(), config.getMaxRetries(), msg.getId()); + } }); } } diff --git a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbMsgDeduplicationNodeTest.java b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbMsgDeduplicationNodeTest.java index 4b9cc0aadd..1a0c3431d7 100644 --- a/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbMsgDeduplicationNodeTest.java +++ b/rule-engine/rule-engine-components/src/test/java/org/thingsboard/rule/engine/transform/TbMsgDeduplicationNodeTest.java @@ -62,11 +62,13 @@ import java.util.function.Consumer; import java.util.stream.Stream; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.ArgumentMatchers.isNull; import static org.mockito.ArgumentMatchers.nullable; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -411,6 +413,80 @@ public class TbMsgDeduplicationNodeTest extends AbstractRuleNodeUpgradeTest { Assertions.assertEquals(msgWithLatestTsInSecondPack.getType(), actualMsg.getType()); } + @Test + public void given_maxRetriesIsZero_when_enqueueFails_then_noRetriesIsScheduled() throws TbNodeException, ExecutionException, InterruptedException { + int wantedNumberOfTellSelfInvocation = 1; + int msgCount = 1; + awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); + invokeTellSelf(wantedNumberOfTellSelfInvocation); + + // Given + when(ctx.getQueueName()).thenReturn(DataConstants.MAIN_QUEUE_NAME); + config.setInterval(deduplicationInterval); + config.setStrategy(DeduplicationStrategy.FIRST); + config.setMaxPendingMsgs(msgCount); + config.setMaxRetries(0); + nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + node.init(ctx, nodeConfiguration); + + DeviceId deviceId = new DeviceId(UUID.randomUUID()); + long currentTimeMillis = System.currentTimeMillis(); + + doAnswer(invocation -> { + Consumer failureCallback = invocation.getArgument(3); + failureCallback.accept(new RuntimeException("Simulated failure")); + return null; + }).when(ctx).enqueueForTellNext(any(), eq(TbNodeConnectionType.SUCCESS), any(), any()); + + TbMsg msg = createMsg(deviceId, currentTimeMillis + 1); + node.onMsg(ctx, msg); + + awaitTellSelfLatch.await(); + + verify(ctx).enqueueForTellNext(any(), eq(TbNodeConnectionType.SUCCESS), any(), any()); + verify(ctx, never()).schedule(any(), anyLong(), any()); + } + + @Test + public void given_maxRetriesIsSetToOne_when_enqueueFails_then_onlyOneRetryIsScheduled() throws TbNodeException, ExecutionException, InterruptedException { + int wantedNumberOfTellSelfInvocation = 1; + int msgCount = 1; + awaitTellSelfLatch = new CountDownLatch(wantedNumberOfTellSelfInvocation); + invokeTellSelf(wantedNumberOfTellSelfInvocation); + + when(ctx.getQueueName()).thenReturn(DataConstants.MAIN_QUEUE_NAME); + config.setInterval(deduplicationInterval); + config.setStrategy(DeduplicationStrategy.FIRST); + config.setMaxPendingMsgs(msgCount); + config.setMaxRetries(1); + nodeConfiguration = new TbNodeConfiguration(JacksonUtil.valueToTree(config)); + node.init(ctx, nodeConfiguration); + + DeviceId deviceId = new DeviceId(UUID.randomUUID()); + long currentTimeMillis = System.currentTimeMillis(); + + doAnswer(invocation -> { + Consumer failureCallback = invocation.getArgument(3); + failureCallback.accept(new RuntimeException("Simulated failure")); + return null; + }).when(ctx).enqueueForTellNext(any(), eq(TbNodeConnectionType.SUCCESS), any(), any()); + + TbMsg msg = createMsg(deviceId, currentTimeMillis + 1); + node.onMsg(ctx, msg); + + awaitTellSelfLatch.await(); + + ArgumentCaptor retryRunnableCaptor = ArgumentCaptor.forClass(Runnable.class); + verify(ctx).schedule(retryRunnableCaptor.capture(), eq(TbMsgDeduplicationNode.TB_MSG_DEDUPLICATION_RETRY_DELAY), eq(TimeUnit.SECONDS)); + + retryRunnableCaptor.getValue().run(); + + // Verify total enqueue attempts (initial + retry) + verify(ctx, times(2)).enqueueForTellNext(any(), eq(TbNodeConnectionType.SUCCESS), any(), any()); + // No more retries scheduled after reaching maxRetries + verify(ctx).schedule(any(), eq(TbMsgDeduplicationNode.TB_MSG_DEDUPLICATION_RETRY_DELAY), eq(TimeUnit.SECONDS)); + } + // Rule nodes upgrade private static Stream givenFromVersionAndConfig_whenUpgrade_thenVerifyHasChangesAndConfig() { return Stream.of( From ebf41af1a231005f3796800934edf8472186bd00 Mon Sep 17 00:00:00 2001 From: yevhenii Date: Tue, 20 May 2025 18:22:24 +0300 Subject: [PATCH 19/28] Update edge install instructions - Minor instruction adjustments --- .../edge/instructions/install/centos/instructions.md | 4 +++- .../edge/instructions/install/docker/instructions.md | 11 +++++++---- .../edge/instructions/install/ubuntu/instructions.md | 6 ++++-- 3 files changed, 14 insertions(+), 7 deletions(-) diff --git a/application/src/main/data/json/edge/instructions/install/centos/instructions.md b/application/src/main/data/json/edge/instructions/install/centos/instructions.md index 90d6d6ddc8..de5e57d926 100644 --- a/application/src/main/data/json/edge/instructions/install/centos/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/centos/instructions.md @@ -41,6 +41,8 @@ OpenJDK 64-Bit Server VM (build ...) #### Step 2. Configure ThingsBoard Database ThingsBoard Edge supports SQL and hybrid database approaches. +In this guide we will use SQL only. +For hybrid details please follow official installation instructions from the ThingsBoard documentation site. ### PostgresSql ThingsBoard Edge uses PostgreSQL database as a local storage. @@ -111,7 +113,7 @@ sudo systemctl restart postgresql-16.service && psql -U postgres -d postgres -h #### Step 3. Choose Queue Service -ThingsBoard Edge is able to use different messaging systems/brokers for storing the messages and communication between ThingsBoard services. +ThingsBoard Edge supports only Kafka or in-memory queue (since v4.0) for message storage and communication between ThingsBoard services. How to choose the right queue implementation? In Memory queue implementation is built-in and default. It is useful for development(PoC) environments and is not suitable for production deployments or any sort of cluster deployments. diff --git a/application/src/main/data/json/edge/instructions/install/docker/instructions.md b/application/src/main/data/json/edge/instructions/install/docker/instructions.md index 3d42cffa8a..be052d9e8a 100644 --- a/application/src/main/data/json/edge/instructions/install/docker/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/docker/instructions.md @@ -12,15 +12,18 @@ Here you can find ThingsBoard Edge docker image: #### Step 2. Choose Queue and/or Database Services -ThingsBoard Edge is able to use different messaging systems/brokers for storing the messages and communication between ThingsBoard services. +ThingsBoard Edge supports only Kafka or in-memory queue (since v4.0) for message storage and communication between ThingsBoard services. + +ThingsBoard Edge supports SQL and hybrid database approaches. +In this guide we will use SQL only. +For hybrid details please follow official installation instructions from the ThingsBoard documentation site. + How to choose the right queue implementation? In Memory queue implementation is built-in and default. It is useful for development(PoC) environments and is not suitable for production deployments or any sort of cluster deployments. Kafka is recommended for production deployments. This queue is used on the most of ThingsBoard production environments now. -Hybrid implementation combines PostgreSQL and Cassandra databases with Kafka queue service. It is recommended if you plan to manage 1M+ devices in production or handle high data ingestion rate (more than 5000 msg/sec). - Create a docker compose file for the ThingsBoard Edge service: ##### In Memory @@ -54,7 +57,7 @@ services: - tb-edge-logs:/var/log/tb-edge postgres: restart: always - image: "postgres:15" + image: "postgres:16" ports: - "5432" environment: diff --git a/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md b/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md index 6ac45d325a..a56b83e558 100644 --- a/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/ubuntu/instructions.md @@ -32,7 +32,9 @@ OpenJDK 64-Bit Server VM (...) #### Step 2. Configure ThingsBoard Edge Database -ThingsBoard Edge supports SQL and hybrid database approaches. See the architecture page for details. +ThingsBoard Edge supports SQL and hybrid database approaches. +In this guide we will use SQL only. +For hybrid details please follow official installation instructions from the ThingsBoard documentation site. ### Configure PostgreSQL ThingsBoard Edge uses PostgreSQL database as a local storage. @@ -71,7 +73,7 @@ echo "CREATE DATABASE tb_edge;" | psql -U postgres -d postgres -h 127.0.0.1 -W #### Step 3. Choose Queue Service -ThingsBoard Edge can use different messaging systems and brokers for storing messages and enabling communication between its services. Choose the appropriate queue implementation based on your specific business needs: +ThingsBoard Edge supports only Kafka or in-memory queue (since v4.0) for message storage and communication between ThingsBoard services. Choose the appropriate queue implementation based on your specific business needs: In Memory: The built-in and default queue implementation. It is useful for development or proof-of-concept (PoC) environments, but is not recommended for production or any type of clustered deployments due to limited scalability. From 7c7886993f95fb9c0264e2ac60b06a6ae3950a2e Mon Sep 17 00:00:00 2001 From: yevhenii Date: Tue, 20 May 2025 18:27:12 +0300 Subject: [PATCH 20/28] Update edge install instructions - Minor instruction adjustments --- .../data/json/edge/instructions/install/docker/instructions.md | 2 ++ 1 file changed, 2 insertions(+) diff --git a/application/src/main/data/json/edge/instructions/install/docker/instructions.md b/application/src/main/data/json/edge/instructions/install/docker/instructions.md index be052d9e8a..e568d273af 100644 --- a/application/src/main/data/json/edge/instructions/install/docker/instructions.md +++ b/application/src/main/data/json/edge/instructions/install/docker/instructions.md @@ -24,6 +24,8 @@ In Memory queue implementation is built-in and default. It is useful for develop Kafka is recommended for production deployments. This queue is used on the most of ThingsBoard production environments now. +Hybrid implementation combines PostgreSQL and Cassandra databases with Kafka queue service. It is recommended if you plan to manage 1M+ devices in production or handle high data ingestion rate (more than 5000 msg/sec). + Create a docker compose file for the ThingsBoard Edge service: ##### In Memory From f990cad134747a831deed49f139801591133af20 Mon Sep 17 00:00:00 2001 From: Yevhenii Date: Wed, 21 May 2025 12:37:42 +0300 Subject: [PATCH 21/28] PostgreSQL View for active Edges - removed testSearch --- .../server/dao/sql/edge/EdgeRepository.java | 11 +++++------ .../thingsboard/server/dao/sql/edge/JpaEdgeDao.java | 5 +---- .../main/resources/sql/schema-views-and-functions.sql | 3 ++- 3 files changed, 8 insertions(+), 11 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java index 66d65cf602..09802e9ca6 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/EdgeRepository.java @@ -44,12 +44,10 @@ public interface EdgeRepository extends JpaRepository { "WHERE d.id = :edgeId") EdgeInfoEntity findEdgeInfoById(@Param("edgeId") UUID edgeId); - @Query(value = "SELECT * " + - "FROM edge_active_attribute_view edge_active " + - "WHERE (:textSearch IS NULL OR edge_active.name ILIKE CONCAT('%', :textSearch, '%')) " + - "ORDER BY edge_active.id", nativeQuery = true) - Page findActiveEdges(@Param("textSearch") String textSearch, - Pageable pageable); + @Query(value = "SELECT * FROM edge_active_attribute_view edge_active", + countQuery = "SELECT count(*) FROM edge_active_attribute_view", + nativeQuery = true) + Page findActiveEdges(Pageable pageable); @Query("SELECT d.id FROM EdgeEntity d WHERE d.tenantId = :tenantId " + "AND (:textSearch IS NULL OR ilike(d.name, CONCAT('%', :textSearch, '%')) = true)") @@ -166,4 +164,5 @@ public interface EdgeRepository extends JpaRepository { @Query("SELECT new org.thingsboard.server.common.data.edqs.fields.EdgeFields(e.id, e.createdTime, e.tenantId, e.customerId," + "e.name, e.version, e.type, e.label, e.additionalInfo) FROM EdgeEntity e WHERE e.id > :id ORDER BY e.id") List findNextBatch(@Param("id") UUID id, Limit limit); + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java index 50b8092731..53ba6cf357 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaEdgeDao.java @@ -68,10 +68,7 @@ public class JpaEdgeDao extends JpaAbstractDao implements Edge @Override public PageData findActiveEdges(PageLink pageLink) { - return DaoUtil.toPageData( - edgeRepository.findActiveEdges( - pageLink.getTextSearch(), - DaoUtil.toPageable(pageLink))); + return DaoUtil.toPageData(edgeRepository.findActiveEdges(DaoUtil.toPageable(pageLink))); } @Override diff --git a/dao/src/main/resources/sql/schema-views-and-functions.sql b/dao/src/main/resources/sql/schema-views-and-functions.sql index b65dc41a90..a6a7ae43cb 100644 --- a/dao/src/main/resources/sql/schema-views-and-functions.sql +++ b/dao/src/main/resources/sql/schema-views-and-functions.sql @@ -89,7 +89,8 @@ SELECT ee.id FROM edge ee JOIN attribute_kv ON ee.id = attribute_kv.entity_id JOIN key_dictionary ON attribute_kv.attribute_key = key_dictionary.key_id -WHERE attribute_kv.bool_v = true AND key_dictionary.key = 'active'; +WHERE attribute_kv.bool_v = true AND key_dictionary.key = 'active' +ORDER BY ee.id; CREATE OR REPLACE FUNCTION create_or_update_active_alarm( t_id uuid, c_id uuid, a_id uuid, a_created_ts bigint, From 48c1842d420650dafbc584635d0831af9e5005b2 Mon Sep 17 00:00:00 2001 From: Yevhenii Date: Wed, 21 May 2025 12:56:23 +0300 Subject: [PATCH 22/28] Add validation for DefaultEdgeRuleChain in DeviceProfile - Refactored parameter name for clarity --- .../dao/service/validator/DeviceProfileDataValidator.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java index e32195022c..dfd0ee82bc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java @@ -209,12 +209,12 @@ public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator Date: Fri, 23 May 2025 15:00:35 +0300 Subject: [PATCH 23/28] Merge pull request #13075 from jekka001/fix-sorting-edge-event Fix sorting edge event --- .../service/edge/rpc/EdgeGrpcSession.java | 2 +- .../rpc/fetch/GeneralEdgeEventFetcher.java | 6 ++- .../server/controller/EdgeControllerTest.java | 48 ++++++++++++------- .../server/edge/AbstractEdgeTest.java | 21 ++++++-- .../server/edge/DeviceProfileEdgeTest.java | 12 +++-- .../server/edge/TenantProfileEdgeTest.java | 9 +++- .../server/edge/imitator/EdgeImitator.java | 15 +++++- .../dao/sql/edge/JpaBaseEdgeEventDao.java | 10 ++-- .../dao/service/EdgeEventServiceTest.java | 25 ++++------ 9 files changed, 98 insertions(+), 50 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java index 7861885989..4a9b68fc6d 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/EdgeGrpcSession.java @@ -712,7 +712,7 @@ public abstract class EdgeGrpcSession implements Closeable { private long findStartSeqIdFromOldestEventIfAny() { long startSeqId = 0L; try { - TimePageLink pageLink = new TimePageLink(1, 0, null, new SortOrder("createdTime"), null, null); + TimePageLink pageLink = new TimePageLink(1, 0, null, null, null, null); PageData edgeEvents = ctx.getEdgeEventService().findEdgeEvents(edge.getTenantId(), edge.getId(), null, null, pageLink); if (!edgeEvents.getData().isEmpty()) { startSeqId = edgeEvents.getData().get(0).getSeqId() - 1; diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java index 0515f23b16..5d7df601b5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/fetch/GeneralEdgeEventFetcher.java @@ -26,9 +26,13 @@ import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.TimePageLink; import org.thingsboard.server.dao.edge.EdgeEventService; +import java.util.concurrent.TimeUnit; + @AllArgsConstructor @Slf4j public class GeneralEdgeEventFetcher implements EdgeEventFetcher { + // Subtract from queueStartTs to ensure no data is lost due to potential misordering of edge events by created_time. + private static final long MISORDERING_COMPENSATION_MILLIS = TimeUnit.SECONDS.toMillis(60); private final Long queueStartTs; private Long seqIdStart; @@ -44,7 +48,7 @@ public class GeneralEdgeEventFetcher implements EdgeEventFetcher { 0, null, null, - queueStartTs, + queueStartTs > 0 ? queueStartTs - MISORDERING_COMPENSATION_MILLIS : 0, System.currentTimeMillis()); } diff --git a/application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java b/application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java index 27c7da09b2..faefab4cb4 100644 --- a/application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/EdgeControllerTest.java @@ -60,6 +60,7 @@ import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.RuleChainId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantProfileId; +import org.thingsboard.server.common.data.id.UserId; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageLink; import org.thingsboard.server.common.data.page.TimePageLink; @@ -68,6 +69,7 @@ import org.thingsboard.server.common.data.rule.RuleChain; import org.thingsboard.server.common.data.rule.RuleChainMetaData; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.DeviceCredentials; +import org.thingsboard.server.common.data.security.UserCredentials; import org.thingsboard.server.common.data.security.model.JwtSettings; import org.thingsboard.server.dao.edge.EdgeDao; import org.thingsboard.server.dao.exception.DataValidationException; @@ -107,6 +109,7 @@ import java.util.concurrent.TimeUnit; import static org.hamcrest.Matchers.containsString; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; +import static org.thingsboard.server.edge.AbstractEdgeTest.CONNECT_MESSAGE_COUNT; @TestPropertySource(properties = { "edges.enabled=true", @@ -138,6 +141,7 @@ public class EdgeControllerTest extends AbstractControllerTest { public EdgeDao edgeDao(EdgeDao edgeDao) { return Mockito.mock(EdgeDao.class, AdditionalAnswers.delegatesTo(edgeDao)); } + } @Before @@ -886,6 +890,8 @@ public class EdgeControllerTest extends AbstractControllerTest { Device savedDevice = doPost("/api/device", device, Device.class); // create public customer + //1 message + // Customer doPost("/api/customer/public/device/" + savedDevice.getId().getId(), Device.class); doDelete("/api/customer/device/" + savedDevice.getId().getId(), Device.class); @@ -897,13 +903,16 @@ public class EdgeControllerTest extends AbstractControllerTest { + "/asset/" + savedAsset.getId().getId().toString(), Asset.class); EdgeImitator edgeImitator = new EdgeImitator(EDGE_HOST, EDGE_PORT, edge.getRoutingKey(), edge.getSecret()); - edgeImitator.ignoreType(UserCredentialsUpdateMsg.class); edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class); edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class); - edgeImitator.expectMessageAmount(27); + // 17 connect message + // + 1 Customer + // + 5 fetchers messages (DeviceProfile, Device, DeviceCredentials, AssetProfile, Asset) in sync process + // + 5 queue messages the same + edgeImitator.expectMessageAmount(CONNECT_MESSAGE_COUNT + 11); edgeImitator.connect(); - waitForMessages(edgeImitator); + edgeImitator.waitForMessages(); verifyFetchersMsgs(edgeImitator, savedDevice); // verify queue msgs @@ -914,9 +923,12 @@ public class EdgeControllerTest extends AbstractControllerTest { Assert.assertTrue(popAssetMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Test Sync Edge Asset 1")); printQueueMsgsIfNotEmpty(edgeImitator); - edgeImitator.expectMessageAmount(21); + // 17 connect messages + // + 1 Customer + // + 5 fetchers messages (DeviceProfile, Device, DeviceCredentials, AssetProfile, Asset) in sync process + edgeImitator.expectMessageAmount(CONNECT_MESSAGE_COUNT + 6); doPost("/api/edge/sync/" + edge.getId()).andExpect(status().isOk()); - waitForMessages(edgeImitator); + edgeImitator.waitForMessages(); verifyFetchersMsgs(edgeImitator, savedDevice); printQueueMsgsIfNotEmpty(edgeImitator); @@ -987,17 +999,6 @@ public class EdgeControllerTest extends AbstractControllerTest { }); } - private void waitForMessages(EdgeImitator edgeImitator) throws Exception { - boolean success = edgeImitator.waitForMessages(); - if (!success) { - List downlinkMsgs = edgeImitator.getDownlinkMsgs(); - for (AbstractMessage downlinkMsg : downlinkMsgs) { - log.error("{}\n{}", downlinkMsg.getClass(), downlinkMsg); - } - Assert.fail("Await for messages was not successful!"); - } - } - private void verifyFetchersMsgs(EdgeImitator edgeImitator, Device savedDevice) { Assert.assertTrue(popQueueMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Main")); Assert.assertTrue(popRuleChainMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Edge Root Rule Chain")); @@ -1011,6 +1012,7 @@ public class EdgeControllerTest extends AbstractControllerTest { Assert.assertTrue(popAssetProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default")); Assert.assertTrue(popDeviceProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default")); Assert.assertTrue(popAssetProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default")); + Assert.assertTrue(popUserCredentialsMsg(edgeImitator.getDownlinkMsgs(), currentUserId)); Assert.assertTrue(popUserMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, TENANT_ADMIN_EMAIL, Authority.TENANT_ADMIN)); Assert.assertTrue(popCustomerMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "Public")); Assert.assertTrue(popDeviceProfileMsg(edgeImitator.getDownlinkMsgs(), UpdateMsgType.ENTITY_CREATED_RPC_MESSAGE, "default")); @@ -1156,6 +1158,20 @@ public class EdgeControllerTest extends AbstractControllerTest { return false; } + private boolean popUserCredentialsMsg(List messages, UserId userId) { + for (AbstractMessage message : messages) { + if (message instanceof UserCredentialsUpdateMsg userCredentialsUpdateMsg) { + UserCredentials userCredentials = JacksonUtil.fromString(userCredentialsUpdateMsg.getEntity(), UserCredentials.class, true); + Assert.assertNotNull(userCredentials); + if (userId.equals(userCredentials.getUserId())) { + messages.remove(message); + return true; + } + } + } + return false; + } + private boolean popUserMsg(List messages, UpdateMsgType msgType, String email, Authority authority) { for (AbstractMessage message : messages) { if (message instanceof UserUpdateMsg userUpdateMsg) { diff --git a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java index d84e0ebea0..c02e72a35f 100644 --- a/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/AbstractEdgeTest.java @@ -115,7 +115,9 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers. }) @Slf4j abstract public class AbstractEdgeTest extends AbstractControllerTest { - + public static final Integer CONNECT_MESSAGE_COUNT = 17; + public static final Integer INSTALLATION_MESSAGE_COUNT = 8; + public static final Integer SYNC_MESSAGE_COUNT = CONNECT_MESSAGE_COUNT + INSTALLATION_MESSAGE_COUNT; private static final String THERMOSTAT_DEVICE_PROFILE_NAME = "Thermostat"; protected DeviceProfile thermostatDeviceProfile; @@ -136,11 +138,12 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { doPost("/api/admin/jwtSettings", settings).andExpect(status().isOk()); loginTenantAdmin(); - + //8 installation messages installation(); edgeImitator = new EdgeImitator("localhost", 7070, edge.getRoutingKey(), edge.getSecret()); - edgeImitator.expectMessageAmount(25); + // 17 connect messages + 8 installation messages + edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT); edgeImitator.ignoreType(OAuth2ClientUpdateMsg.class); edgeImitator.ignoreType(OAuth2DomainUpdateMsg.class); edgeImitator.connect(); @@ -164,22 +167,32 @@ abstract public class AbstractEdgeTest extends AbstractControllerTest { thermostatDeviceProfile = this.createDeviceProfile(THERMOSTAT_DEVICE_PROFILE_NAME, createMqttDeviceProfileTransportConfiguration(new JsonTransportPayloadConfiguration(), false)); extendDeviceProfileData(thermostatDeviceProfile); + //2 messages DeviceProfile thermostatDeviceProfile = doPost("/api/deviceProfile", thermostatDeviceProfile, DeviceProfile.class); Device savedDevice = saveDevice("Edge Device 1", THERMOSTAT_DEVICE_PROFILE_NAME); // create public customer + //1 message + // Customer doPost("/api/customer/public/device/" + savedDevice.getId().getId(), Device.class); doDelete("/api/customer/device/" + savedDevice.getId().getId(), Device.class); - Asset savedAsset = saveAsset("Edge Asset 1"); + Asset savedAsset = saveAsset("Edge Asset 1"); updateRootRuleChainMetadata(); edge = doPost("/api/edge", constructEdge("Test Edge", "test"), Edge.class); + //3 messages + // Device + // DeviceProfile + // DeviceCredentials doPost("/api/edge/" + edge.getUuidId() + "/device/" + savedDevice.getUuidId(), Device.class); + //2 messages + // Asset + // AssetProfile doPost("/api/edge/" + edge.getUuidId() + "/asset/" + savedAsset.getUuidId(), Asset.class); diff --git a/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java index 819dec5c0d..a29d7a4672 100644 --- a/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/DeviceProfileEdgeTest.java @@ -128,15 +128,19 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { @Test public void testDeleteDeviceProfilesWhenEdgeIsOffline() throws Exception { + //2 message RuleChain and RuleChainMetadata RuleChainId thermostatsRuleChainId = createEdgeRuleChainAndAssignToEdge("Thermostats Rule Chain"); // create device profile DeviceProfile deviceProfile = this.createDeviceProfile("ONE_MORE_DEVICE_PROFILE", null); deviceProfile.setDefaultEdgeRuleChainId(thermostatsRuleChainId); extendDeviceProfileData(deviceProfile); + + //1 message DeviceProfile edgeImitator.expectMessageAmount(1); deviceProfile = doPost("/api/deviceProfile", deviceProfile, DeviceProfile.class); Assert.assertTrue(edgeImitator.waitForMessages()); + AbstractMessage latestMessage = edgeImitator.getLatestMessage(); Assert.assertTrue(latestMessage instanceof DeviceProfileUpdateMsg); DeviceProfileUpdateMsg deviceProfileUpdateMsg = (DeviceProfileUpdateMsg) latestMessage; @@ -150,9 +154,11 @@ public class DeviceProfileEdgeTest extends AbstractEdgeTest { doDelete("/api/deviceProfile/" + deviceProfile.getUuidId()) .andExpect(status().isOk()); edgeImitator.connect(); - // 27 sync message - // + 1 delete message - edgeImitator.expectMessageAmount(28); + + // 25 sync message + // + 2 RuleChain and RuleChainMetadata + // + 1 delete DeviceProfile + edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT + 3); Assert.assertTrue(edgeImitator.waitForMessages()); latestMessage = edgeImitator.getLatestMessage(); diff --git a/application/src/test/java/org/thingsboard/server/edge/TenantProfileEdgeTest.java b/application/src/test/java/org/thingsboard/server/edge/TenantProfileEdgeTest.java index 51cfa44ff9..c7e9ed1575 100644 --- a/application/src/test/java/org/thingsboard/server/edge/TenantProfileEdgeTest.java +++ b/application/src/test/java/org/thingsboard/server/edge/TenantProfileEdgeTest.java @@ -78,11 +78,14 @@ public class TenantProfileEdgeTest extends AbstractEdgeTest { TenantProfileQueueConfiguration mainQueueConfiguration = createQueueConfig(DataConstants.MAIN_QUEUE_NAME, DataConstants.MAIN_QUEUE_TOPIC); TenantProfileQueueConfiguration isolatedQueueConfiguration = createQueueConfig("IsolatedHighPriority", "tb_rule_engine.isolated_hp"); edgeTenantProfile.getProfileData().setQueueConfiguration(List.of(mainQueueConfiguration, isolatedQueueConfiguration)); + // + 1 TenantProfile + // + 1 Queue main + // + 1 Queue isolated edgeImitator.expectMessageAmount(3); edgeTenantProfile = doPost("/api/tenantProfile", edgeTenantProfile, TenantProfile.class); Assert.assertTrue(edgeImitator.waitForMessages()); - Optional tenantProfileUpdateMsgOpt = edgeImitator.findMessageByType(TenantProfileUpdateMsg.class); + Optional tenantProfileUpdateMsgOpt = edgeImitator.findMessageByType(TenantProfileUpdateMsg.class); Assert.assertTrue(tenantProfileUpdateMsgOpt.isPresent()); TenantProfileUpdateMsg tenantProfileUpdateMsg = tenantProfileUpdateMsgOpt.get(); TenantProfile tenantProfile = JacksonUtil.fromString(tenantProfileUpdateMsg.getEntity(), TenantProfile.class, true); @@ -96,7 +99,9 @@ public class TenantProfileEdgeTest extends AbstractEdgeTest { loginTenantAdmin(); - edgeImitator.expectMessageAmount(21); + // 25 sync message + // +1 isolated Queue + edgeImitator.expectMessageAmount(SYNC_MESSAGE_COUNT + 1); doPost("/api/edge/sync/" + edge.getId()); assertThat(edgeImitator.waitForMessages()).as("await for messages after edge sync rest api call").isTrue(); diff --git a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java index be9492bc53..8b259cf8fc 100644 --- a/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java +++ b/application/src/test/java/org/thingsboard/server/edge/imitator/EdgeImitator.java @@ -24,6 +24,7 @@ import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.checkerframework.checker.nullness.qual.Nullable; +import org.junit.Assert; import org.thingsboard.edge.rpc.EdgeGrpcClient; import org.thingsboard.edge.rpc.EdgeRpcClient; import org.thingsboard.server.controller.AbstractWebTest; @@ -386,7 +387,19 @@ public class EdgeImitator { } public boolean waitForMessages() throws InterruptedException { - return waitForMessages(AbstractWebTest.TIMEOUT); + boolean success = waitForMessages(AbstractWebTest.TIMEOUT); + + if (!success) { + List downlinkMsgs = getDownlinkMsgs(); + for (AbstractMessage downlinkMsg : downlinkMsgs) { + log.error("{}\n{}", downlinkMsg.getClass(), downlinkMsg); + } + + log.error("message count: {}", downlinkMsgs.size()); + Assert.fail("Await for messages was not successful!"); + } + + return true; } public boolean waitForMessages(int timeoutInSeconds) throws InterruptedException { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java index 4b6b898f42..ea2abbba8a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/edge/JpaBaseEdgeEventDao.java @@ -44,7 +44,7 @@ import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository; import org.thingsboard.server.dao.util.SqlDao; -import java.util.ArrayList; +import java.util.Collections; import java.util.Comparator; import java.util.List; import java.util.UUID; @@ -58,6 +58,7 @@ import static org.thingsboard.server.dao.model.ModelConstants.NULL_UUID; @RequiredArgsConstructor @Slf4j public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao implements EdgeEventDao { + private static final List SORT_ORDERS = Collections.singletonList(new SortOrder("seqId")); private final UUID systemTenantId = NULL_UUID; @@ -175,11 +176,6 @@ public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao findEdgeEvents(UUID tenantId, EdgeId edgeId, Long seqIdStart, Long seqIdEnd, TimePageLink pageLink) { - List sortOrders = new ArrayList<>(); - if (pageLink.getSortOrder() != null) { - sortOrders.add(pageLink.getSortOrder()); - } - sortOrders.add(new SortOrder("seqId")); return DaoUtil.toPageData( edgeEventRepository .findEdgeEventsByTenantIdAndEdgeId( @@ -190,7 +186,7 @@ public class JpaBaseEdgeEventDao extends JpaPartitionedAbstractDao> futures = new ArrayList<>(); - futures.add(saveEdgeEventWithProvidedTime(timeBeforeStartTime, edgeId, deviceId, tenantId)); - futures.add(saveEdgeEventWithProvidedTime(eventTime, edgeId, deviceId, tenantId)); - futures.add(saveEdgeEventWithProvidedTime(eventTime + 1, edgeId, deviceId, tenantId)); - futures.add(saveEdgeEventWithProvidedTime(eventTime + 2, edgeId, deviceId, tenantId)); - futures.add(saveEdgeEventWithProvidedTime(timeAfterEndTime, edgeId, deviceId, tenantId)); - - Futures.allAsList(futures).get(); + saveEdgeEventWithProvidedTime(timeBeforeStartTime, edgeId, deviceId, tenantId).get(); + saveEdgeEventWithProvidedTime(eventTime, edgeId, deviceId, tenantId).get(); + saveEdgeEventWithProvidedTime(eventTime + 2, edgeId, deviceId, tenantId).get(); + saveEdgeEventWithProvidedTime(eventTime + 1, edgeId, deviceId, tenantId).get(); + saveEdgeEventWithProvidedTime(timeAfterEndTime, edgeId, deviceId, tenantId).get(); TimePageLink pageLink = new TimePageLink(2, 0, "", new SortOrder("createdTime", SortOrder.Direction.DESC), startTime, endTime); PageData edgeEvents = edgeEventService.findEdgeEvents(tenantId, edgeId, 0L, null, pageLink); Assert.assertNotNull(edgeEvents.getData()); Assert.assertEquals(2, edgeEvents.getData().size()); - Assert.assertEquals(Uuids.startOf(eventTime + 2), edgeEvents.getData().get(0).getUuidId()); - Assert.assertEquals(Uuids.startOf(eventTime + 1), edgeEvents.getData().get(1).getUuidId()); + Assert.assertEquals(Uuids.startOf(eventTime), edgeEvents.getData().get(0).getUuidId()); + Assert.assertEquals(Uuids.startOf(eventTime + 2), edgeEvents.getData().get(1).getUuidId()); Assert.assertTrue(edgeEvents.hasNext()); Assert.assertNotNull(pageLink.nextPageLink()); @@ -130,7 +124,7 @@ public class EdgeEventServiceTest extends AbstractServiceTest { Assert.assertNotNull(edgeEvents.getData()); Assert.assertEquals(1, edgeEvents.getData().size()); - Assert.assertEquals(Uuids.startOf(eventTime), edgeEvents.getData().get(0).getUuidId()); + Assert.assertEquals(Uuids.startOf(eventTime + 1), edgeEvents.getData().get(0).getUuidId()); Assert.assertFalse(edgeEvents.hasNext()); edgeEventDao.cleanupEvents(1); @@ -141,4 +135,5 @@ public class EdgeEventServiceTest extends AbstractServiceTest { edgeEvent.setId(new EdgeEventId(Uuids.startOf(time))); return edgeEventService.saveAsync(edgeEvent); } + } \ No newline at end of file From ae4bdb52da5b82056431e63662661b08f06fcfb5 Mon Sep 17 00:00:00 2001 From: Yevhenii Date: Fri, 23 May 2025 18:02:19 +0300 Subject: [PATCH 24/28] Add validation for DefaultEdgeRuleChain in DeviceProfile - refactoring --- .../dao/service/validator/DeviceProfileDataValidator.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java index dfd0ee82bc..ca312608d7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java @@ -189,11 +189,11 @@ public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator Date: Fri, 23 May 2025 18:10:33 +0300 Subject: [PATCH 25/28] Add validation for DefaultEdgeRuleChain in DeviceProfile - reverted changes --- .../dao/service/validator/DeviceProfileDataValidator.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java index ca312608d7..dfd0ee82bc 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/DeviceProfileDataValidator.java @@ -189,11 +189,11 @@ public class DeviceProfileDataValidator extends AbstractHasOtaPackageValidator Date: Tue, 27 May 2025 17:23:25 +0300 Subject: [PATCH 26/28] Add Serializable to TrendzSettings --- .../thingsboard/server/common/data/trendz/TrendzSettings.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/trendz/TrendzSettings.java b/common/data/src/main/java/org/thingsboard/server/common/data/trendz/TrendzSettings.java index 3c3b49399c..ab637efff7 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/trendz/TrendzSettings.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/trendz/TrendzSettings.java @@ -17,8 +17,10 @@ package org.thingsboard.server.common.data.trendz; import lombok.Data; +import java.io.Serializable; + @Data -public class TrendzSettings { +public class TrendzSettings implements Serializable { private boolean enabled; private String baseUrl; From 80a33492d587c932452acd3584ee82aec04e1cf6 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Tue, 27 May 2025 18:26:19 +0300 Subject: [PATCH 27/28] Add proto number for ComponentLifecycleEvent --- .../server/common/data/EntityType.java | 17 ++++++++ .../data/plugin/ComponentLifecycleEvent.java | 40 ++++++++++++++++--- .../server/common/util/ProtoUtils.java | 21 +++++----- .../server/common/util/ProtoUtilsTest.java | 7 ++++ 4 files changed, 69 insertions(+), 16 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java index 93e754eb2c..7582b377f0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/EntityType.java @@ -18,6 +18,7 @@ package org.thingsboard.server.common.data; import lombok.Getter; import org.apache.commons.lang3.StringUtils; +import java.util.Arrays; import java.util.EnumSet; import java.util.List; @@ -76,6 +77,15 @@ public enum EntityType { public static final List NORMAL_NAMES = EnumSet.allOf(EntityType.class).stream() .map(EntityType::getNormalName).toList(); + private static final EntityType[] BY_PROTO; + + static { + BY_PROTO = new EntityType[Arrays.stream(values()).mapToInt(EntityType::getProtoNumber).max().orElse(0) + 1]; + for (EntityType entityType : values()) { + BY_PROTO[entityType.getProtoNumber()] = entityType; + } + } + EntityType(int protoNumber) { this.protoNumber = protoNumber; this.tableName = name().toLowerCase(); @@ -98,4 +108,11 @@ public enum EntityType { return false; } + public static EntityType forProtoNumber(int protoNumber) { + if (protoNumber < 0 || protoNumber >= BY_PROTO.length) { + throw new IllegalArgumentException("Invalid EntityType proto number " + protoNumber); + } + return BY_PROTO[protoNumber]; + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/plugin/ComponentLifecycleEvent.java b/common/data/src/main/java/org/thingsboard/server/common/data/plugin/ComponentLifecycleEvent.java index 32ced72292..5d13db2348 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/plugin/ComponentLifecycleEvent.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/plugin/ComponentLifecycleEvent.java @@ -15,12 +15,42 @@ */ package org.thingsboard.server.common.data.plugin; +import lombok.Getter; +import lombok.RequiredArgsConstructor; + import java.io.Serializable; +import java.util.Arrays; -/** - * @author Andrew Shvayka - */ +@RequiredArgsConstructor public enum ComponentLifecycleEvent implements Serializable { - // In sync with ComponentLifecycleEvent proto - CREATED, STARTED, ACTIVATED, SUSPENDED, UPDATED, STOPPED, DELETED, FAILED, DEACTIVATED + + CREATED(0), + STARTED(1), + ACTIVATED(2), + SUSPENDED(3), + UPDATED(4), + STOPPED(5), + DELETED(6), + FAILED(7), + DEACTIVATED(8); + + @Getter + private final int protoNumber; // corresponds to ComponentLifecycleEvent proto + + private static final ComponentLifecycleEvent[] BY_PROTO; + + static { + BY_PROTO = new ComponentLifecycleEvent[Arrays.stream(values()).mapToInt(ComponentLifecycleEvent::getProtoNumber).max().orElse(0) + 1]; + for (ComponentLifecycleEvent event : values()) { + BY_PROTO[event.getProtoNumber()] = event; + } + } + + public static ComponentLifecycleEvent forProtoNumber(int protoNumber) { + if (protoNumber < 0 || protoNumber >= BY_PROTO.length) { + throw new IllegalArgumentException("Invalid ComponentLifecycleEvent proto number " + protoNumber); + } + return BY_PROTO[protoNumber]; + } + } diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java index 502884eb31..77250486de 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java +++ b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java @@ -102,7 +102,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.ApiUsageRecordKeyPro import org.thingsboard.server.gen.transport.TransportProtos.KeyValueProto; import java.util.ArrayList; -import java.util.Arrays; import java.util.List; import java.util.Set; import java.util.UUID; @@ -114,14 +113,6 @@ import static org.thingsboard.server.common.data.DataConstants.GATEWAY_PARAMETER @Slf4j public class ProtoUtils { - private static final EntityType[] entityTypeByProtoNumber; - - static { - int arraySize = Arrays.stream(EntityType.values()).mapToInt(EntityType::getProtoNumber).max().orElse(0); - entityTypeByProtoNumber = new EntityType[arraySize + 1]; - Arrays.stream(EntityType.values()).forEach(entityType -> entityTypeByProtoNumber[entityType.getProtoNumber()] = entityType); - } - public static TransportProtos.ComponentLifecycleMsgProto toProto(ComponentLifecycleMsg msg) { var builder = TransportProtos.ComponentLifecycleMsgProto.newBuilder() .setTenantIdMSB(msg.getTenantId().getId().getMostSignificantBits()) @@ -129,7 +120,7 @@ public class ProtoUtils { .setEntityType(toProto(msg.getEntityId().getEntityType())) .setEntityIdMSB(msg.getEntityId().getId().getMostSignificantBits()) .setEntityIdLSB(msg.getEntityId().getId().getLeastSignificantBits()) - .setEvent(TransportProtos.ComponentLifecycleEvent.forNumber(msg.getEvent().ordinal())); + .setEvent(TransportProtos.ComponentLifecycleEvent.forNumber(msg.getEvent().getProtoNumber())); if (msg.getProfileId() != null) { builder.setProfileIdMSB(msg.getProfileId().getId().getMostSignificantBits()); builder.setProfileIdLSB(msg.getProfileId().getId().getLeastSignificantBits()); @@ -175,7 +166,15 @@ public class ProtoUtils { } public static EntityType fromProto(TransportProtos.EntityTypeProto entityType) { - return entityTypeByProtoNumber[entityType.getNumber()]; + return EntityType.forProtoNumber(entityType.getNumber()); + } + + public static TransportProtos.ComponentLifecycleEvent toProto(ComponentLifecycleEvent event) { + return TransportProtos.ComponentLifecycleEvent.forNumber(event.getProtoNumber()); + } + + public static ComponentLifecycleEvent fromProto(TransportProtos.ComponentLifecycleEvent eventProto) { + return ComponentLifecycleEvent.forProtoNumber(eventProto.getNumber()); } public static TransportProtos.ToEdgeSyncRequestMsgProto toProto(ToEdgeSyncRequest request) { diff --git a/common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java b/common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java index e1714a9335..25788721dc 100644 --- a/common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java +++ b/common/proto/src/test/java/org/thingsboard/server/common/util/ProtoUtilsTest.java @@ -114,6 +114,13 @@ class ProtoUtilsTest { } } + @Test + void protoComponentLifecycleEventSerialization() { + for (ComponentLifecycleEvent event : ComponentLifecycleEvent.values()) { + assertThat(ProtoUtils.fromProto(ProtoUtils.toProto(event))).isEqualTo(event); + } + } + @Test void protoEdgeHighPrioritySerialization() { EdgeHighPriorityMsg msg = new EdgeHighPriorityMsg(tenantId, EdgeUtils.constructEdgeEvent(tenantId, edgeId, From 731e8df5acfa6c7952a202a5111bef20147ab8fa Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Wed, 28 May 2025 11:08:01 +0300 Subject: [PATCH 28/28] Fix ComponentLifecycleMsg toProto --- .../java/org/thingsboard/server/common/util/ProtoUtils.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java index 77250486de..bb2e0fb2d1 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java +++ b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java @@ -120,7 +120,7 @@ public class ProtoUtils { .setEntityType(toProto(msg.getEntityId().getEntityType())) .setEntityIdMSB(msg.getEntityId().getId().getMostSignificantBits()) .setEntityIdLSB(msg.getEntityId().getId().getLeastSignificantBits()) - .setEvent(TransportProtos.ComponentLifecycleEvent.forNumber(msg.getEvent().getProtoNumber())); + .setEvent(toProto(msg.getEvent())); if (msg.getProfileId() != null) { builder.setProfileIdMSB(msg.getProfileId().getId().getMostSignificantBits()); builder.setProfileIdLSB(msg.getProfileId().getId().getLeastSignificantBits()); @@ -147,7 +147,7 @@ public class ProtoUtils { var builder = ComponentLifecycleMsg.builder() .tenantId(TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB()))) .entityId(entityId) - .event(ComponentLifecycleEvent.values()[proto.getEventValue()]); + .event(fromProto(proto.getEvent())); if (!StringUtils.isEmpty(proto.getName())) { builder.name(proto.getName()); }