diff --git a/application/pom.xml b/application/pom.xml
index dba7d27c1b..4cbd9c3b64 100644
--- a/application/pom.xml
+++ b/application/pom.xml
@@ -419,6 +419,10 @@
+
+ dev.langchain4j
+ langchain4j-ollama
+
diff --git a/application/src/main/data/json/system/widget_types/markdown_html_card.json b/application/src/main/data/json/system/widget_types/markdown_html_card.json
index 64a847952a..b66e3ee88e 100644
--- a/application/src/main/data/json/system/widget_types/markdown_html_card.json
+++ b/application/src/main/data/json/system/widget_types/markdown_html_card.json
@@ -11,16 +11,10 @@
"resources": [],
"templateHtml": "\n",
"templateCss": "#container tb-markdown-widget {\n height: 100%;\n display: block;\n}\n\n#container tb-markdown-widget .tb-markdown-view {\n height: 100%;\n overflow: auto;\n}\n",
- "controllerScript": "self.onInit = function() {\n}\n\nself.onDataUpdated = function() {\n self.ctx.$scope.markdownWidget.onDataUpdated();\n}\n\nself.actionSources = function() {\n return {\n 'elementClick': {\n name: 'widget-action.element-click',\n multiple: true\n }\n };\n}\n\nself.typeParameters = function() {\n return {\n dataKeysOptional: true,\n datasourcesOptional: true,\n hasDataPageLink: true\n };\n}\n\nself.onDestroy = function() {\n}\n\n",
- "settingsSchema": "",
- "dataKeySettingsSchema": "",
+ "controllerScript": "self.onInit = function() {\n}\n\nself.onDataUpdated = function() {\n self.ctx.$scope.markdownWidget.onDataUpdated();\n}\n\nself.actionSources = function() {\n return {\n 'elementClick': {\n name: 'widget-action.element-click',\n multiple: true\n }\n };\n}\n\nself.typeParameters = function() {\n return {\n dataKeysOptional: true,\n datasourcesOptional: true,\n hasDataPageLink: true,\n hideDataSettings: true\n };\n}\n\nself.onDestroy = function() {\n}\n\n",
"settingsDirective": "tb-markdown-widget-settings",
- "defaultConfig": "{\"datasources\":[{\"type\":\"function\",\"name\":\"function\",\"dataKeys\":[{\"name\":\"f(x)\",\"type\":\"function\",\"label\":\"Random\",\"color\":\"#2196f3\",\"settings\":{},\"_hash\":0.15479322438769105,\"funcBody\":\"var value = prevValue + Math.random() * 100 - 50;\\nvar multiplier = Math.pow(10, 2 || 0);\\nvar value = Math.round(value * multiplier) / multiplier;\\nif (value < -1000) {\\n\\tvalue = -1000;\\n} else if (value > 1000) {\\n\\tvalue = 1000;\\n}\\nreturn value;\"}]}],\"timewindow\":{\"realtime\":{\"timewindowMs\":60000}},\"showTitle\":false,\"backgroundColor\":\"#fff\",\"color\":\"rgba(0, 0, 0, 0.87)\",\"padding\":\"0px\",\"settings\":{\"markdownTextPattern\":\"### Markdown/HTML card\\n - **Current entity**: ${entityName}.\\n - **Current value**: ${Random}.\",\"markdownTextFunction\":\"return '# Some title\\\\n - Entity name: ' + data[0]['entityName'];\",\"useMarkdownTextFunction\":false},\"title\":\"Markdown/HTML Card\",\"showTitleIcon\":false,\"iconColor\":\"rgba(0, 0, 0, 0.87)\",\"iconSize\":\"24px\",\"titleTooltip\":\"\",\"dropShadow\":true,\"enableFullscreen\":true,\"widgetStyle\":{},\"titleStyle\":{\"fontSize\":\"16px\",\"fontWeight\":400},\"showLegend\":false}"
+ "defaultConfig": "{\"datasources\":[{\"type\":\"function\",\"name\":\"function\",\"dataKeys\":[{\"name\":\"f(x)\",\"type\":\"function\",\"label\":\"Temperature\",\"color\":\"#2196f3\",\"settings\":{},\"_hash\":0.15479322438769105,\"funcBody\":\"const baseTemp = 20;\\nconst dailySwing = 10;\\nconst hourlyVariation = Math.sin((time % 24) * Math.PI / 12) * dailySwing;\\nconst randomness = (Math.random() - 0.5) * 2;\\nconst smoothingFactor = 0.8;\\nreturn (prevValue * smoothingFactor) + ((baseTemp + hourlyVariation + randomness) * (1 - smoothingFactor));\",\"aggregationType\":null,\"units\":null,\"decimals\":null,\"usePostProcessing\":null,\"postFuncBody\":null}],\"alarmFilterConfig\":{\"statusList\":[\"ACTIVE\"]}}],\"timewindow\":{\"realtime\":{\"timewindowMs\":60000}},\"showTitle\":false,\"backgroundColor\":\"#fff\",\"color\":\"rgba(0, 0, 0, 0.87)\",\"padding\":\"0px\",\"settings\":{\"useMarkdownTextFunction\":false,\"markdownTextPattern\":\"### Markdown/HTML card\\n - **Current entity**: ${entityName}.\\n - **Current value**: ${Temperature}.\",\"markdownTextFunction\":\"return '# Some title\\\\n - Entity name: ' + data[0]['entityName'];\",\"applyDefaultMarkdownStyle\":true,\"markdownCss\":\"\"},\"title\":\"Markdown/HTML Card\",\"showTitleIcon\":false,\"iconColor\":\"rgba(0, 0, 0, 0.87)\",\"iconSize\":\"24px\",\"titleTooltip\":\"\",\"dropShadow\":true,\"enableFullscreen\":true,\"widgetStyle\":{},\"titleStyle\":{\"fontSize\":\"16px\",\"fontWeight\":400},\"showLegend\":false,\"useDashboardTimewindow\":true,\"displayTimewindow\":true}"
},
- "tags": [
- "web",
- "markup"
- ],
"resources": [
{
"link": "/api/images/system/markdown_html_card_system_widget_image.png",
@@ -33,5 +27,10 @@
"data": "iVBORw0KGgoAAAANSUhEUgAAAMgAAACgCAMAAAB+IdObAAABdFBMVEX////u7u7g4ODf398nJydISEiCs/Tx8fE/Pz+amppgYGC6urr7+/upqan6+vrW1tbKysr19fX9/f3j4+M5OTn39/fs7OxSUlJxcXHFxcVCQkI8PDyMjIxFRUV+fn40NDTa2trp6el3d3dVVVXAwMCKiorQ0NClpaX6/P8vLy+SkpJbW1srKyvl5eWhoaHn5+dMTEytra2Hh4dlZWVPT0+3t7eenp6Pj49CjO7j7v3MzMx7e3uBgYE2NjabwvbR0dHIyMjV5vvY2NgyMjJra2tUl/A0hO3a6fzA2fqEhIQ4hu1XV1f3+v7Dw8OXl5eUlJT0+P6y0Pi0tLTc3NzV1dWysrKwsLCnyfeIt/U+ie6vr690dHRJSUkqfezy9/641Pnz8/NkoPJIj++rq6vg7fyszfh8r/Tt7e28vLySvfZZmvDR4/uNuvV2q/Pp8v1ppPJMku9oaGifxfctf+zT09NnZ2dfnvHL3/vG3fpup/Lt9P4vgO3lU88CAAAQR0lEQVR42uyby2/aQBDGJ7jINti8DOb9dAA1lEASHhIQJWp7QSiJ1JYcaA6VckivkdL/v4aPLbs1bkqLlZTmOwzKesa7P8WKP80QIq1s+P5xGeWQzWGYez+L/jFpsqFR2dz750GI5DIZe7sAohnk2wkQ8r2ArBQ6o6XCUXpEMS7TS5DmLQcyjaVNIS6UmBFdC8eQ/LSUX3psu6K0NjP2absg5538CiSRp6QpxIVKr3zGq7d/CJJQE2sz5fhWQcz7+gokeJxX3wX5iIqLRu+wkaaj1HBAdH6c+tyyQUKXTaoNC3GJSoVRifQy0Z0u51KFMEFYCedyBtEy8+0oVZXJN6pEyBi/zuV823y0HlYgtUrl0qrxERXX6sm9ek1vS8HTL3RyUIskJH/ivUrNbPBrUTKVYFcxP5+3moOUpNSOsiFUYcVMdyK0zGxmu2eKZOaDEcWQg8V0WvYIhIb7el+I0GG98rneIznYzZfoJE1EktI4IqoX5g/MPA71/mxU+dCX/KtHCCtEgQjLPFxEPdDtBo7waHkEcnU6ydxc8REVB7pp6ld7mWovG16CZK0DIrU6P9g8VtXZQaBxNZ2DJMOowgpAlpnHi9hQVbXmKYisyFY4ykdUHM+IZtXIhMhiIP7EtxxF2lH7YOfJUKgTKSnqpVLiQbBC9D5Iy8wzJfwpLkU6i9tqrxNbBUkpbzIXS5DwSUyJCREqTImmhUSgHSiWGAiFrBnlFOudQceZTJW0otkqajwIVoiC8c7+MvM68L5o0FXWurHXh5nk/hO9EOVyiESFFivRKP1K0btlZixEzdO5+8Z9ZPnp3+x/pmYmf/PheVmUZ6EdBQkW+h93ASQcOCp0dgHEVvmNuRsg9cBu/EaCmfBOgKSVwU781bq7+bobf37DrzKZzHQHQHbnhfgC8vRyAdlqLwuKkSBPQczZqC+tQD49IG7WyxJrhY4WL09ByuPbyoSBaL7DnqEhbtbLEmuFjpYoLx+tj0UGIrUtqy0hbtjLEmqxz7KjNaikrgjC3bwC0U8uGQiVG4EW4sa9LL4WQkcrEb/4eE4Q7uYVyEHg3mQgh7p+iLhxL4uvhdAIit1UDWLC3bx6tO7a+wBh+vNeFiSAkKz6RwThbt6A2L+MVubcCbJ5L8shdLRiIdKyPvyMu3kDEslPlJHwHvmLXpYodLTKmUb7RMPPuJtHj1brobX43GIvS+xoaS1ZuNvTv9mfnf4HkNbt7YDM5dwiJrgsp9xz3N2apyCRDz9A9i19QGeFbER0SpLfZernmuMUcrwEOcsGViAp9vcfTok7gHPq55LzVCDliboWBE4JUz/ptFJ4u5z6wV9hDuiWg30Gw1x/JDvdWvmQqLcnV1Nj3zZBqr2va0HglDD1k7KRo6yGqR/8FeaAbjlLb3A1HI9Up1sLJ+feQFLObrPmFpsPlrkeBBFTv/kjkffhsYG/whzQLYeZnPpld+xwawsQ5Ff07YEcZxqd00tXEEz98O7GIeGvMAd05ogg1/Xu2OnWSh2Wn1O3ByI9POgTSQCBR0LE1A+HxNQP/gpzQLccBkLUHTvdmlz8MgfJx7RkeksgkPhoMY+EiKkfDompH/wV5oBuOTzIGrfWVyw7/40/M/TohQgQeCRENvWDuKkf5oAuOY+6tYSp2VBRmTwDuYhbqnPq58lkUPK/eK1HQeCyNtHg9raFp4fmmkZ/x00hE1H0YLHNHJoIct7r9VoAgcvaCES39oloL9ldWE4lxjsuRKeQiYgcVLn0wXD1cZDxZwbC3h0bKbU/Dz20Lo55x4XoFDIRkYMqOLc/Bxle20EAgTtCtwo+Cu4I69x3rhhIqY1nohEh3nEhwlOx7lYwyDIRkYMqODd+X1RJ8dGwRuik4QxrQe5TelkEgTtCtwo+CqYC6/x3rpYgl2jAmVmNeMeFCE/FulsTi2UiIgdVcG78vqiSil+D2RY6aTjD+tGbfp8sCyBwR+hWwUcBBOvCd64AMpwSETpevONChKdi3a1WE5mIyGFV2J3fd16FlUIdnTScwQkCtfcFELgjvKnho+COsC585wogYx21F0So4kHgqfjuFjIRkYMq7M7vi6r5SlVFJw1nWAticn0teCe4I9wQPgruCOvCd64AclTBsUNEvONChKdi3a16HZmIyGFV2J3fF1X2SqIdQScNZ1gHUspwfS14J7gj3BA+Cu4I68J3rgAiZw2ycT7bgXdciPBUrLsVaCCTReSgCrvz+6LKeNdQcoROGs7wWF8L3klwR/BRcEdYF75zBRC6tWSi+6BgsGREeCp0t9ANQybiKpPtLuyLKjTSUIszbG5R4KPcV1QrfrH41Cl6Gv09T4VMxF/7txevxSRMCWMOj4SrLg6q/IVc5TGIGTx0gvglZ+cK7ghXXR3ULOluO70FCXfe99xA4H9Ej+QOgqv9CrnKU5BJffVoofuEKaEwB4TPgUfCVazIl6mKAUfEHJQGSo/lBKkV+/0BA0H3CVNCfg4InwN3xM8QE+1cLS3DEeGqrQOVvJcTZHZzVI+XfoCkiU0J+fEZfA4++RniRYOITQlxxdYUH95LBJme2AZYF0DU6s8g8DlwR/wMcZoiYlNCXLWlV8l7OUHCSuuu81UAwZQQvgj+Bz4H7oifIfriBrEpIa7aGl6T5xJBoFy+k9oTQDAlhC+C/4HPgTvCVbaSTbbZlBBXqalEyVO5/48VfnKdEkbvBI8kXo21ovBR7Gp0sr1fiCcWBe7o8at6j9z1TECenV5Anpt2CeQ7u3b0mjYQwHH8l/kTE03aNUZtrVFrO0Ps1TnbKjilYvIioxaq60v7IPRBXwvr/79LRre6MdhgGemW70vCwR33gQsEkn8kSfknSiBxK4HErQQStxJI3EogcSuBxK3/DTLR8njqvPcby28hSNEG+K6IIVbNCfddq13gefe8xFOP7q+vzj6C5mzh+ybzKCGCJci6ZPNPQPRZ/acQ9UOkkBEngG5ZAWQ7/OicN1LOpYTImz2YudnrEDK+ck6AG8m73AOObnD/eqI5O5AVjRRM4xV0Y5AxLoCTdnsQQF71nKKxDSjObIyCMbKNQoSQW68LrFmVkDbv0mwjx31OJWRlW69wS8v1JKTKsueNcWrDVPdNlE/hqp7KLGRL0UWTQ2zTCI5WfSRUT0KanlAtaphbns1K0RaurUQJeW/p8NOOhOxWTbNjS8hpannPXMlb4IINM2O7GHMKxS6bPb79RC4OpNctK/CZgez4EH2W4fAkgOyKBQwJybopOBKSHR2gK/JRHy1/wLUitJ6EoNmu2J6EXAP3LPMcaHMQPiNvKGe1mJrwulXqXJ2zDvcxGJlDdiV2Og1OdtOQkHDDc7ZM9wFYUNtiyTCmzEUNyeLOb4viTEKmwn8GEWIXmLEeQobMfFGpw1KrX5qqeA5p8h0ndzN1GkJqjyFki40QsqJVkv0FiCbSPlps7ogGcPoV0uvzBjfMyTE3NKDKAzRULi5YbmxA9JFaRsXmOoSc2SbqbGH/EPhEDbVjADqihxQF1wFky+2Me8J9glzqh7V8xkuvg7Fi7WzPCLb+kappqrzegMDnFE2KZQjR2D3KSsgb9uVVk9fKQLNXOLOOCpFC8GDpAQQfPd5lefIEwdz9oB9Z7GRdYM8mfQVQOASGXG1CDI5hqscIIYVbIaoSsvQpHiSk0HXpXQGGYD4yyGZ6ET9krvAlZYlfbVkIZxYzhWuef/tLeJlB2Et7aXSs/sxSNzb/MiHLSid9m8fXXixkswSSQF5KCSRuJZC4lUDiVgL53N699jQNxXEc/6HHTWVDGXL1Ak6Ezg0F5dISLRTsWmvFFZWuKqut28LuczO7wJv3dJTBiGgMooh8k2U77U6yT5om/ydLT1udEO/u3zOXbs1hv14yg936iR9HNt6N7xRNMnEVu/FpHFFdgxsr/BZIHyErAJ6T4K9DhsePgGQOQfQ8DpczD0Gy0nEhq4+BrsgqhYz51/sxMHfB76WQt3O9uO5/40CWJ4cDWP6KT3P0bamve+BzdwC0lWvAxOR4AE6P3s94XUjAZKBqWlGmEEEXkCvUGskcUN9GTAYkvQE466Rel5Jgc2YqwZgKzx0TErzhw8ORcBCTkScjsz3D5NnqAJn5unoL8+TeZQp5Ebq6OoQ7YWyQCdyc95PpMLkF2q3HmJy+QneCgm6v3Rnpda/IZgy1Wl5kwbOFfDQnW6JWF8EoFcgWtDKfzaNSQE6JFwwerBgvZRleiUvHhHSvbmHw9b0gNrawSN4Ok/foJQs3PgR8o4/paXpR1unx5e5Qz1Bo4QIZ85MLGBp0Ic8n0fWgBZkAPl5vQxosgwaFlIA0Dz5NFYmGWEO8wigSElGVQpo8nBcrQ6VyVsBxIcHnK2RqOoi+4PRlcn2YvKGQCLmIKbLl3CMbpAtYveKNLI9emX1/M+An/XhxyYUsjoQnP4HWvxAeDI23IXoTEFhHgHim9dY0PWZ5W0wmq5Zl7UgUUq47ypagWvwNkPXrkQ9XcTuIl0/uruxBwrNhXxeZcyDjZAo9kYd4efVSb+hJEB0QBCaC90BbCPZjcB+ilQ5DUh5bzVdsFHeSHMcxFCI2fjMkcJusO5DLQwMLZHkXMjMWWsN0eMpP/N5nT7vWIhfxgNzHF/KuEzL7ClMRH4ChBxgb3WpDOIVDpQ2hvxrFaBM5xQOUU2B0B5IvQBXbELt+bAjWRvsdyEaIfCGTLgTzkYmlm+RqyI/hy+TZDLBIFjFDBjohr27P3ngI2tLNkcHZ+TYEFcW22hAhWgLEFMDKgFQuGy1bosba+5C4oh0D0lnPJ3QW8KJVlw9H5vV1ftctwahmFnupsY5zKpwYNYGahr22mdM1omQ8rRQxG7U8P8xSsmUlvb/mTxfELaFVJPwkSU/FzofG/wsS22TwneJ5nKbOKCQhcykZgt4AwKU0FaqcTDkQoQ5niDXB5QRdaEOSKXo00Sim5NaGhiRB3gZMZ6Mute7gnMABgl7HCdcJkaKlDGuVMoYG2eYLIrgd0UMhdUPA5jaKVWhKM6+YLkQrV5oWBKWUsXlIiidtZGDngCinimneqKOh5JuGDr3MizwOdbIQhYHGMqgU0JAQiHLczjZim3Ujhz1IFtBEFyInoVYTghKDVkMhAxT2IEUNyHhQ00FfapRDMcqgs5OFGIAsAqkSEvGsWJW4KCjEKKANsQAu6kK2PVlxsyiwgJyFKAPxPQhMK8umYUhAUxd2LMuqcujsT0EKeQaGC0mx5kGIZLiQEs+gugepmQcgOZsDn6YLB5JUONrfuiLZFOo7wi6EyRkcjDpkCrFjiDddSNmEvMm5kEyJUbMZ1FKQqpwpMqpVQNqDhK0ztgZVx6H+FERWWIs1XQh4kUkpdsmBiHaZcyGmwpaUhguJWYYtZug+u0Zv9prBNrMoimy5rEOgI+6fvNk7YxJo1Tm0ahYOHGdiBz8nAiW9vS/BOFvoQpTd1cnWCfH8NIs96kyJjq/RdOexmpEtG+5nE0f0d2YtyfzR+JrAoYSKpuKHnQ+N55BT3jnktHUOOW2dIUifD2cgXx/uenEG8t5FT5/3n78mzkO0z8hjzX34BpPgTEZLPbVzAAAAAElFTkSuQmCC",
"public": true
}
+ ],
+ "scada": false,
+ "tags": [
+ "web",
+ "markup"
]
}
\ No newline at end of file
diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql
index 016e786776..0add4c0545 100644
--- a/application/src/main/data/upgrade/basic/schema_update.sql
+++ b/application/src/main/data/upgrade/basic/schema_update.sql
@@ -14,3 +14,35 @@
-- limitations under the License.
--
+-- UPDATE TENANT PROFILE CONFIGURATION START
+
+UPDATE tenant_profile
+SET profile_data = jsonb_set(
+ profile_data,
+ '{configuration}',
+ (profile_data -> 'configuration')
+ || jsonb_strip_nulls(
+ jsonb_build_object(
+ 'minAllowedScheduledUpdateIntervalInSecForCF',
+ CASE
+ WHEN (profile_data -> 'configuration') ? 'minAllowedScheduledUpdateIntervalInSecForCF'
+ THEN NULL
+ ELSE to_jsonb(60)
+ END,
+ 'maxRelationLevelPerCfArgument',
+ CASE
+ WHEN (profile_data -> 'configuration') ? 'maxRelationLevelPerCfArgument'
+ THEN NULL
+ ELSE to_jsonb(10)
+ END
+ )
+ ),
+ false
+ )
+WHERE NOT (
+ (profile_data -> 'configuration') ? 'minAllowedScheduledUpdateIntervalInSecForCF'
+ AND
+ (profile_data -> 'configuration') ? 'maxRelationLevelPerCfArgument'
+ );
+
+-- UPDATE TENANT PROFILE CONFIGURATION END
diff --git a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
index ea46ce86eb..b23015a1fe 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ActorSystemContext.java
@@ -97,6 +97,7 @@ import org.thingsboard.server.dao.ota.OtaPackageService;
import org.thingsboard.server.dao.queue.QueueService;
import org.thingsboard.server.dao.queue.QueueStatsService;
import org.thingsboard.server.dao.relation.RelationService;
+import org.thingsboard.server.dao.resource.TbResourceDataCache;
import org.thingsboard.server.dao.resource.ResourceService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.rule.RuleNodeStateService;
@@ -511,6 +512,10 @@ public class ActorSystemContext {
@Getter
private ResourceService resourceService;
+ @Autowired
+ @Getter
+ private TbResourceDataCache resourceDataCache;
+
@Lazy
@Autowired(required = false)
@Getter
diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java
index 2959bfc8eb..c57984ef3d 100644
--- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityActor.java
@@ -51,7 +51,7 @@ public class CalculatedFieldEntityActor extends AbstractCalculatedFieldActor {
@Override
public void destroy(TbActorStopReason stopReason, Throwable cause) throws TbActorException {
log.debug("[{}] Stopping CF entity actor.", processor.tenantId);
- processor.stop();
+ processor.stop(false);
}
@Override
diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
index 35539834c3..f8b61a082f 100644
--- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
@@ -48,6 +48,8 @@ import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
+import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry;
+import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
import java.util.ArrayList;
import java.util.Collection;
@@ -91,16 +93,18 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
this.ctx = ctx;
}
- public void stop() {
- log.info("[{}][{}] Stopping entity actor.", tenantId, entityId);
+ public void stop(boolean partitionChanged) {
+ log.info(partitionChanged ?
+ "[{}][{}] Stopping entity actor due to change partition event." :
+ "[{}][{}] Stopping entity actor.",
+ tenantId, entityId);
states.clear();
ctx.stop(ctx.getSelf());
}
public void process(CalculatedFieldPartitionChangeMsg msg) {
if (!systemContext.getPartitionService().resolve(ServiceType.TB_RULE_ENGINE, DataConstants.CF_QUEUE_NAME, tenantId, entityId).isMyPartition()) {
- log.info("[{}] Stopping entity actor due to change partition event.", entityId);
- ctx.stop(ctx.getSelf());
+ stop(true);
}
}
@@ -251,6 +255,16 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
if (state == null) {
state = getOrInitState(ctx);
justRestored = true;
+ } else if (ctx.shouldFetchDynamicArgumentsFromDb(state)) {
+ log.debug("[{}][{}] Going to update dynamic arguments for CF.", entityId, ctx.getCfId());
+ try {
+ Map dynamicArgsFromDb = cfService.fetchDynamicArgsFromDb(ctx, entityId);
+ dynamicArgsFromDb.forEach(newArgValues::putIfAbsent);
+ var geofencingState = (GeofencingCalculatedFieldState) state;
+ geofencingState.setLastDynamicArgumentsRefreshTs(System.currentTimeMillis());
+ } catch (Exception e) {
+ throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build();
+ }
}
if (state.isSizeOk()) {
if (state.updateState(ctx, newArgValues) || justRestored) {
@@ -271,7 +285,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
if (state != null) {
return state;
} else {
- ListenableFuture stateFuture = systemContext.getCalculatedFieldProcessingService().fetchStateFromDb(ctx, entityId);
+ ListenableFuture stateFuture = cfService.fetchStateFromDb(ctx, entityId);
// Ugly but necessary. We do not expect to often fetch data from DB. Only once per pair lifetime.
// This call happens while processing the CF pack from the queue consumer. So the timeout should be relatively low.
// Alternatively, we can fetch the state outside the actor system and push separate command to create this actor,
@@ -288,7 +302,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
boolean stateSizeChecked = false;
try {
if (ctx.isInitialized() && state.isReady()) {
- CalculatedFieldResult calculationResult = state.performCalculation(ctx).get(systemContext.getCfCalculationResultTimeout(), TimeUnit.SECONDS);
+ CalculatedFieldResult calculationResult = state.performCalculation(entityId, ctx).get(systemContext.getCfCalculationResultTimeout(), TimeUnit.SECONDS);
state.checkStateSize(ctxId, ctx.getMaxStateSize());
stateSizeChecked = true;
if (state.isSizeOk()) {
@@ -298,7 +312,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
callback.onSuccess();
}
if (DebugModeUtil.isDebugAllAvailable(ctx.getCalculatedField())) {
- systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, calculationResult.getResult().toString(), null);
+ systemContext.persistCalculatedFieldDebugEvent(tenantId, ctx.getCfId(), entityId, state.getArguments(), tbMsgId, tbMsgType, calculationResult.toStringOrElseNull(), null);
}
}
} else {
@@ -367,7 +381,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
}
private Map mapToArguments(CalculatedFieldCtx ctx, AttributeScopeProto scope, List attrDataList) {
- return mapToArguments(ctx.getMainEntityArguments(), scope, attrDataList);
+ return mapToArguments(entityId, ctx.getMainEntityArguments(), ctx.getMainEntityGeofencingArgumentNames(), scope, attrDataList);
}
private Map mapToArguments(CalculatedFieldCtx ctx, EntityId entityId, AttributeScopeProto scope, List attrDataList) {
@@ -375,17 +389,23 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
if (argNames.isEmpty()) {
return Collections.emptyMap();
}
- return mapToArguments(argNames, scope, attrDataList);
+ List geofencingArgumentNames = ctx.getLinkedEntityGeofencingArgumentNames();
+ return mapToArguments(entityId, argNames, geofencingArgumentNames, scope, attrDataList);
}
- private Map mapToArguments(Map argNames, AttributeScopeProto scope, List attrDataList) {
+ private Map mapToArguments(EntityId entityId, Map argNames, List geofencingArgNames, AttributeScopeProto scope, List attrDataList) {
Map arguments = new HashMap<>();
for (AttributeValueProto item : attrDataList) {
ReferencedEntityKey key = new ReferencedEntityKey(item.getKey(), ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name()));
String argName = argNames.get(key);
- if (argName != null) {
- arguments.put(argName, new SingleValueArgumentEntry(item));
+ if (argName == null) {
+ continue;
+ }
+ if (geofencingArgNames.contains(argName)) {
+ arguments.put(argName, new GeofencingArgumentEntry(entityId, item));
+ continue;
}
+ arguments.put(argName, new SingleValueArgumentEntry(item));
}
return arguments;
}
@@ -395,26 +415,32 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
if (argNames.isEmpty()) {
return Collections.emptyMap();
}
- return mapToArgumentsWithDefaultValue(argNames, ctx.getArguments(), scope, removedAttrKeys);
+ List geofencingArgumentNames = ctx.getLinkedEntityGeofencingArgumentNames();
+ return mapToArgumentsWithDefaultValue(argNames, ctx.getArguments(), geofencingArgumentNames, scope, removedAttrKeys);
}
private Map mapToArgumentsWithDefaultValue(CalculatedFieldCtx ctx, AttributeScopeProto scope, List removedAttrKeys) {
- return mapToArgumentsWithDefaultValue(ctx.getMainEntityArguments(), ctx.getArguments(), scope, removedAttrKeys);
+ return mapToArgumentsWithDefaultValue(ctx.getMainEntityArguments(), ctx.getArguments(), ctx.getMainEntityGeofencingArgumentNames(), scope, removedAttrKeys);
}
- private Map mapToArgumentsWithDefaultValue(Map argNames, Map configArguments, AttributeScopeProto scope, List removedAttrKeys) {
+ private Map mapToArgumentsWithDefaultValue(Map argNames, Map configArguments, List geofencingArgNames, AttributeScopeProto scope, List removedAttrKeys) {
Map arguments = new HashMap<>();
for (String removedKey : removedAttrKeys) {
ReferencedEntityKey key = new ReferencedEntityKey(removedKey, ArgumentType.ATTRIBUTE, AttributeScope.valueOf(scope.name()));
String argName = argNames.get(key);
- if (argName != null) {
- Argument argument = configArguments.get(argName);
- String defaultValue = (argument != null) ? argument.getDefaultValue() : null;
- arguments.put(argName, StringUtils.isNotEmpty(defaultValue)
- ? new SingleValueArgumentEntry(System.currentTimeMillis(), new StringDataEntry(removedKey, defaultValue), null)
- : new SingleValueArgumentEntry());
-
+ if (argName == null) {
+ continue;
}
+ if (geofencingArgNames.contains(argName)) {
+ arguments.put(argName, new GeofencingArgumentEntry());
+ continue;
+ }
+ Argument argument = configArguments.get(argName);
+ String defaultValue = (argument != null) ? argument.getDefaultValue() : null;
+ arguments.put(argName, StringUtils.isNotEmpty(defaultValue)
+ ? new SingleValueArgumentEntry(System.currentTimeMillis(), new StringDataEntry(removedKey, defaultValue), null)
+ : new SingleValueArgumentEntry());
+
}
return arguments;
}
diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
index c00995d3d4..7a76cb9821 100644
--- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
+++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
@@ -57,6 +57,7 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.function.BiConsumer;
import static org.thingsboard.server.utils.CalculatedFieldUtils.fromProto;
@@ -255,7 +256,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
log.debug("[{}] Failed to lookup CF by id [{}]", tenantId, cfId);
callback.onSuccess();
} else {
- var cfCtx = new CalculatedFieldCtx(cf, systemContext.getTbelInvokeService(), systemContext.getApiLimitService());
+ var cfCtx = getCfCtx(cf);
try {
cfCtx.init();
} catch (Exception e) {
@@ -266,11 +267,15 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
// Alternative approach would be to use any list but avoid modifications to the list (change the complete map value instead)
entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cfCtx);
addLinks(cf);
- initCf(cfCtx, callback, false);
+ applyToTargetCfEntityActors(cfCtx, callback, (id, cb) -> initCfForEntity(id, cfCtx, false, cb));
}
}
}
+ private CalculatedFieldCtx getCfCtx(CalculatedField cf) {
+ return new CalculatedFieldCtx(cf, systemContext.getTbelInvokeService(), systemContext.getApiLimitService(), systemContext.getRelationService());
+ }
+
private void onCfUpdated(ComponentLifecycleMsg msg, TbCallback callback) throws CalculatedFieldException {
var cfId = new CalculatedFieldId(msg.getEntityId().getId());
var oldCfCtx = calculatedFields.get(cfId);
@@ -282,7 +287,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
log.debug("[{}] Failed to lookup CF by id [{}]", tenantId, cfId);
callback.onSuccess();
} else {
- var newCfCtx = new CalculatedFieldCtx(newCf, systemContext.getTbelInvokeService(), systemContext.getApiLimitService());
+ var newCfCtx = getCfCtx(newCf);
try {
newCfCtx.init();
} catch (Exception e) {
@@ -290,6 +295,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
}
calculatedFields.put(newCf.getId(), newCfCtx);
List oldCfList = entityIdCalculatedFields.get(newCf.getEntityId());
+
List newCfList = new CopyOnWriteArrayList<>();
boolean found = false;
for (CalculatedFieldCtx oldCtx : oldCfList) {
@@ -312,7 +318,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
// Alternative approach would be to use any list but avoid modifications to the list (change the complete map value instead)
var stateChanges = newCfCtx.hasStateChanges(oldCfCtx);
if (stateChanges || newCfCtx.hasOtherSignificantChanges(oldCfCtx)) {
- initCf(newCfCtx, callback, stateChanges);
+ applyToTargetCfEntityActors(newCfCtx, callback, (id, cb) -> initCfForEntity(id, newCfCtx, stateChanges, cb));
} else {
callback.onSuccess();
}
@@ -326,31 +332,11 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
if (cfCtx == null) {
log.debug("[{}] CF was already deleted [{}]", tenantId, cfId);
callback.onSuccess();
- } else {
- entityIdCalculatedFields.get(cfCtx.getEntityId()).remove(cfCtx);
- deleteLinks(cfCtx);
-
- EntityId entityId = cfCtx.getEntityId();
- EntityType entityType = cfCtx.getEntityId().getEntityType();
- if (isProfileEntity(entityType)) {
- var entityIds = entityProfileCache.getEntityIdsByProfileId(entityId);
- if (!entityIds.isEmpty()) {
- //TODO: no need to do this if we cache all created actors and know which one belong to us;
- var multiCallback = new MultipleTbCallback(entityIds.size(), callback);
- entityIds.forEach(id -> {
- if (isMyPartition(id, multiCallback)) {
- deleteCfForEntity(id, cfId, multiCallback);
- }
- });
- } else {
- callback.onSuccess();
- }
- } else {
- if (isMyPartition(entityId, callback)) {
- deleteCfForEntity(entityId, cfId, callback);
- }
- }
+ return;
}
+ entityIdCalculatedFields.get(cfCtx.getEntityId()).remove(cfCtx);
+ deleteLinks(cfCtx);
+ applyToTargetCfEntityActors(cfCtx, callback, (id, cb) -> deleteCfForEntity(id, cfId, cb));
}
public void onTelemetryMsg(CalculatedFieldTelemetryMsg msg) {
@@ -389,31 +375,11 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
}
for (var linkProto : linksList) {
var link = fromProto(linkProto);
- var targetEntityId = link.entityId();
- var targetEntityType = targetEntityId.getEntityType();
var cf = calculatedFields.get(link.cfId());
- if (EntityType.DEVICE_PROFILE.equals(targetEntityType) || EntityType.ASSET_PROFILE.equals(targetEntityType)) {
- // iterate over all entities that belong to profile and push the message for corresponding CF
- var entityIds = entityProfileCache.getEntityIdsByProfileId(targetEntityId);
- if (!entityIds.isEmpty()) {
- MultipleTbCallback multipleCallback = new MultipleTbCallback(entityIds.size(), callback);
- var newMsg = new EntityCalculatedFieldLinkedTelemetryMsg(tenantId, sourceEntityId, proto.getMsg(), cf, multipleCallback);
- entityIds.forEach(entityId -> {
- if (isMyPartition(entityId, multipleCallback)) {
- log.debug("Pushing linked telemetry msg to specific actor [{}]", entityId);
- getOrCreateActor(entityId).tell(newMsg);
- }
- });
- } else {
- callback.onSuccess();
- }
- } else {
- if (isMyPartition(targetEntityId, callback)) {
- log.debug("Pushing linked telemetry msg to specific actor [{}]", targetEntityId);
- var newMsg = new EntityCalculatedFieldLinkedTelemetryMsg(tenantId, sourceEntityId, proto.getMsg(), cf, callback);
- getOrCreateActor(targetEntityId).tell(newMsg);
- }
- }
+ withTargetEntities(link.entityId(), callback, (ids, cb) -> {
+ var linkedTelemetryMsg = new EntityCalculatedFieldLinkedTelemetryMsg(tenantId, sourceEntityId, proto.getMsg(), cf, cb);
+ ids.forEach(id -> linkedTelemetryMsgForEntity(id, linkedTelemetryMsg));
+ });
}
}
@@ -452,26 +418,9 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
return result;
}
- private void initCf(CalculatedFieldCtx cfCtx, TbCallback callback, boolean forceStateReinit) {
- EntityId entityId = cfCtx.getEntityId();
- EntityType entityType = cfCtx.getEntityId().getEntityType();
- if (isProfileEntity(entityType)) {
- var entityIds = entityProfileCache.getEntityIdsByProfileId(entityId);
- if (!entityIds.isEmpty()) {
- var multiCallback = new MultipleTbCallback(entityIds.size(), callback);
- entityIds.forEach(id -> {
- if (isMyPartition(id, multiCallback)) {
- initCfForEntity(id, cfCtx, forceStateReinit, multiCallback);
- }
- });
- } else {
- callback.onSuccess();
- }
- } else {
- if (isMyPartition(entityId, callback)) {
- initCfForEntity(entityId, cfCtx, forceStateReinit, callback);
- }
- }
+ private void linkedTelemetryMsgForEntity(EntityId entityId, EntityCalculatedFieldLinkedTelemetryMsg msg) {
+ log.debug("Pushing linked telemetry msg to specific actor [{}]", entityId);
+ getOrCreateActor(entityId).tell(msg);
}
private void deleteCfForEntity(EntityId entityId, CalculatedFieldId cfId, TbCallback callback) {
@@ -545,7 +494,7 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
}
private void initCalculatedField(CalculatedField cf) throws CalculatedFieldException {
- var cfCtx = new CalculatedFieldCtx(cf, systemContext.getTbelInvokeService(), systemContext.getApiLimitService());
+ var cfCtx = new CalculatedFieldCtx(cf, systemContext.getTbelInvokeService(), systemContext.getApiLimitService(), systemContext.getRelationService());
try {
cfCtx.init();
} catch (Exception e) {
@@ -584,4 +533,30 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
}
}
+ private void applyToTargetCfEntityActors(CalculatedFieldCtx ctx,
+ TbCallback callback,
+ BiConsumer action) {
+ withTargetEntities(ctx.getEntityId(), callback, (ids, cb) -> ids.forEach(id -> action.accept(id, cb)));
+ }
+
+ private void withTargetEntities(EntityId entityId, TbCallback parentCallback, BiConsumer, TbCallback> consumer) {
+ if (isProfileEntity(entityId.getEntityType())) {
+ var ids = entityProfileCache.getEntityIdsByProfileId(entityId);
+ if (ids.isEmpty()) {
+ parentCallback.onSuccess();
+ return;
+ }
+ var multiCallback = new MultipleTbCallback(ids.size(), parentCallback);
+ var profileEntityIds = ids.stream().filter(id -> isMyPartition(id, multiCallback)).toList();
+ if (profileEntityIds.isEmpty()) {
+ return;
+ }
+ consumer.accept(profileEntityIds, multiCallback);
+ return;
+ }
+ if (isMyPartition(entityId, parentCallback)) {
+ consumer.accept(List.of(entityId), parentCallback);
+ }
+ }
+
}
diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/MultipleTbCallback.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/MultipleTbCallback.java
index d1f4c9092e..493985c97a 100644
--- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/MultipleTbCallback.java
+++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/MultipleTbCallback.java
@@ -50,7 +50,7 @@ public class MultipleTbCallback implements TbCallback {
@Override
public void onFailure(Throwable t) {
- log.warn("[{}][{}] onFailure.", id, callback.getId());
+ log.warn("[{}][{}] onFailure.", id, callback.getId(), t);
callback.onFailure(t);
}
}
diff --git a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
index 6374e4016d..88b04c7613 100644
--- a/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
+++ b/application/src/main/java/org/thingsboard/server/actors/ruleChain/DefaultTbContext.java
@@ -51,6 +51,7 @@ import org.thingsboard.server.common.data.Device;
import org.thingsboard.server.common.data.DeviceProfile;
import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.HasRuleEngineProfile;
+import org.thingsboard.server.common.data.HasTenantId;
import org.thingsboard.server.common.data.StringUtils;
import org.thingsboard.server.common.data.TenantProfile;
import org.thingsboard.server.common.data.alarm.Alarm;
@@ -60,6 +61,7 @@ import org.thingsboard.server.common.data.id.AssetId;
import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.DeviceId;
import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.HasId;
import org.thingsboard.server.common.data.id.RuleChainId;
import org.thingsboard.server.common.data.id.RuleNodeId;
import org.thingsboard.server.common.data.id.TenantId;
@@ -110,6 +112,7 @@ import org.thingsboard.server.dao.ota.OtaPackageService;
import org.thingsboard.server.dao.queue.QueueService;
import org.thingsboard.server.dao.queue.QueueStatsService;
import org.thingsboard.server.dao.relation.RelationService;
+import org.thingsboard.server.dao.resource.TbResourceDataCache;
import org.thingsboard.server.dao.resource.ResourceService;
import org.thingsboard.server.dao.rule.RuleChainService;
import org.thingsboard.server.dao.tenant.TenantService;
@@ -770,6 +773,11 @@ public class DefaultTbContext implements TbContext {
return mainCtx.getResourceService();
}
+ @Override
+ public TbResourceDataCache getTbResourceDataCache() {
+ return mainCtx.getResourceDataCache();
+ }
+
@Override
public OtaPackageService getOtaPackageService() {
return mainCtx.getOtaPackageService();
@@ -1054,7 +1062,18 @@ public class DefaultTbContext implements TbContext {
@Override
public void checkTenantEntity(EntityId entityId) throws TbNodeException {
- if (!this.getTenantId().equals(TenantIdLoader.findTenantId(this, entityId))) {
+ TenantId actualTenantId = TenantIdLoader.findTenantId(this, entityId);
+ assertSameTenantId(actualTenantId, entityId);
+ }
+
+ @Override
+ public & HasTenantId, I extends EntityId> void checkTenantEntity(E entity) throws TbNodeException {
+ TenantId actualTenantId = entity.getTenantId();
+ assertSameTenantId(actualTenantId, entity.getId());
+ }
+
+ private void assertSameTenantId(TenantId tenantId, EntityId entityId) throws TbNodeException {
+ if (!getTenantId().equals(tenantId)) {
throw new TbNodeException("Entity with id: '" + entityId + "' specified in the configuration doesn't belong to the current tenant.", true);
}
}
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 5945355ef8..c5b077c128 100644
--- a/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/CalculatedFieldController.java
@@ -289,7 +289,6 @@ public class CalculatedFieldController extends BaseController {
default -> throw new IllegalArgumentException("Calculated fields do not support '" + entityType + "' for referenced entities.");
}
}
-
}
}
diff --git a/application/src/main/java/org/thingsboard/server/controller/ImageController.java b/application/src/main/java/org/thingsboard/server/controller/ImageController.java
index f9ec7fd844..9288484f86 100644
--- a/application/src/main/java/org/thingsboard/server/controller/ImageController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/ImageController.java
@@ -300,6 +300,7 @@ public class ImageController extends BaseController {
tbImageService.putETag(cacheKey, descriptor.getEtag());
var result = ResponseEntity.ok()
.header("Content-Type", descriptor.getMediaType())
+ .header("Content-Security-Policy", "default-src 'none'")
.eTag(descriptor.getEtag());
if (!cacheKey.isPublic()) {
result
diff --git a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java
index 29f4daa783..b9968aefa9 100644
--- a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java
@@ -162,6 +162,8 @@ public class SystemInfoController extends BaseController {
}
systemParams.setMaxArgumentsPerCF(tenantProfileConfiguration.getMaxArgumentsPerCF());
systemParams.setMaxDataPointsPerRollingArg(tenantProfileConfiguration.getMaxDataPointsPerRollingArg());
+ systemParams.setMinAllowedScheduledUpdateIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedScheduledUpdateIntervalInSecForCF());
+ systemParams.setMaxRelationLevelPerCfArgument(tenantProfileConfiguration.getMaxRelationLevelPerCfArgument());
systemParams.setTrendzSettings(trendzSettingsService.findTrendzSettings(currentUser.getTenantId()));
}
systemParams.setMobileQrEnabled(Optional.ofNullable(qrCodeSettingService.findQrCodeSettings(TenantId.SYS_TENANT_ID))
diff --git a/application/src/main/java/org/thingsboard/server/controller/TbResourceController.java b/application/src/main/java/org/thingsboard/server/controller/TbResourceController.java
index b23603f6a2..54d2494679 100644
--- a/application/src/main/java/org/thingsboard/server/controller/TbResourceController.java
+++ b/application/src/main/java/org/thingsboard/server/controller/TbResourceController.java
@@ -55,14 +55,17 @@ import org.thingsboard.server.common.data.util.ThrowingSupplier;
import org.thingsboard.server.config.annotations.ApiOperation;
import org.thingsboard.server.queue.util.TbCoreComponent;
import org.thingsboard.server.service.resource.TbResourceService;
+import org.thingsboard.server.service.security.model.SecurityUser;
import org.thingsboard.server.service.security.permission.Operation;
import org.thingsboard.server.service.security.permission.Resource;
+import java.util.ArrayList;
import java.util.Collections;
import java.util.EnumSet;
import java.util.HashSet;
import java.util.List;
import java.util.Set;
+import java.util.UUID;
import static org.thingsboard.server.controller.ControllerConstants.AVAILABLE_FOR_ANY_AUTHORIZED_USER;
import static org.thingsboard.server.controller.ControllerConstants.LWM2M_OBJECT_DESCRIPTION;
@@ -263,6 +266,20 @@ public class TbResourceController extends BaseController {
}
}
+ @ApiOperation(value = "Get Resource Infos by ids (getSystemOrTenantResourcesByIds)")
+ @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN')")
+ @GetMapping(value = "/resource", params = {"resourceIds"})
+ public List getSystemOrTenantResourcesByIds(
+ @Parameter(description = "A list of resource ids, separated by comma ','", array = @ArraySchema(schema = @Schema(type = "string")))
+ @RequestParam("resourceIds") Set resourceUuids) throws ThingsboardException {
+ SecurityUser user = getCurrentUser();
+ List resourceIds = new ArrayList<>();
+ for (UUID resourceId : resourceUuids) {
+ resourceIds.add(new TbResourceId(resourceId));
+ }
+ return resourceService.findSystemOrTenantResourcesByIds(user.getTenantId(), resourceIds);
+ }
+
@ApiOperation(value = "Get All Resource Infos (getAllResources)",
notes = "Returns a page of Resource Info objects owned by tenant. " +
PAGE_DATA_PARAMETERS + RESOURCE_INFO_DESCRIPTION + TENANT_AUTHORITY_PARAGRAPH)
diff --git a/application/src/main/java/org/thingsboard/server/service/ai/Langchain4jChatModelConfigurerImpl.java b/application/src/main/java/org/thingsboard/server/service/ai/Langchain4jChatModelConfigurerImpl.java
index 69dd98f47f..c631f1a3d0 100644
--- a/application/src/main/java/org/thingsboard/server/service/ai/Langchain4jChatModelConfigurerImpl.java
+++ b/application/src/main/java/org/thingsboard/server/service/ai/Langchain4jChatModelConfigurerImpl.java
@@ -32,8 +32,10 @@ import dev.langchain4j.model.chat.request.ChatRequestParameters;
import dev.langchain4j.model.github.GitHubModelsChatModel;
import dev.langchain4j.model.googleai.GoogleAiGeminiChatModel;
import dev.langchain4j.model.mistralai.MistralAiChatModel;
+import dev.langchain4j.model.ollama.OllamaChatModel;
import dev.langchain4j.model.openai.OpenAiChatModel;
import dev.langchain4j.model.vertexai.gemini.VertexAiGeminiChatModel;
+import org.springframework.http.HttpHeaders;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.ai.model.chat.AmazonBedrockChatModelConfig;
import org.thingsboard.server.common.data.ai.model.chat.AnthropicChatModelConfig;
@@ -43,10 +45,12 @@ import org.thingsboard.server.common.data.ai.model.chat.GoogleAiGeminiChatModelC
import org.thingsboard.server.common.data.ai.model.chat.GoogleVertexAiGeminiChatModelConfig;
import org.thingsboard.server.common.data.ai.model.chat.Langchain4jChatModelConfigurer;
import org.thingsboard.server.common.data.ai.model.chat.MistralAiChatModelConfig;
+import org.thingsboard.server.common.data.ai.model.chat.OllamaChatModelConfig;
import org.thingsboard.server.common.data.ai.model.chat.OpenAiChatModelConfig;
import org.thingsboard.server.common.data.ai.provider.AmazonBedrockProviderConfig;
import org.thingsboard.server.common.data.ai.provider.AzureOpenAiProviderConfig;
import org.thingsboard.server.common.data.ai.provider.GoogleVertexAiGeminiProviderConfig;
+import org.thingsboard.server.common.data.ai.provider.OllamaProviderConfig;
import software.amazon.awssdk.auth.credentials.AwsBasicCredentials;
import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider;
import software.amazon.awssdk.regions.Region;
@@ -54,7 +58,11 @@ import software.amazon.awssdk.services.bedrockruntime.BedrockRuntimeClient;
import java.io.ByteArrayInputStream;
import java.io.IOException;
+import java.nio.charset.StandardCharsets;
import java.time.Duration;
+import java.util.Base64;
+
+import static java.util.Collections.singletonMap;
@Component
class Langchain4jChatModelConfigurerImpl implements Langchain4jChatModelConfigurer {
@@ -62,6 +70,7 @@ class Langchain4jChatModelConfigurerImpl implements Langchain4jChatModelConfigur
@Override
public ChatModel configureChatModel(OpenAiChatModelConfig chatModelConfig) {
return OpenAiChatModel.builder()
+ .baseUrl(chatModelConfig.providerConfig().baseUrl())
.apiKey(chatModelConfig.providerConfig().apiKey())
.modelName(chatModelConfig.modelId())
.temperature(chatModelConfig.temperature())
@@ -134,7 +143,7 @@ class Langchain4jChatModelConfigurerImpl implements Langchain4jChatModelConfigur
// set request timeout from model config
if (chatModelConfig.timeoutSeconds() != null) {
- retrySettings.setTotalTimeout(org.threeten.bp.Duration.ofSeconds(chatModelConfig.timeoutSeconds()));
+ retrySettings.setTotalTimeoutDuration(Duration.ofSeconds(chatModelConfig.timeoutSeconds()));
}
// set updated retry settings
@@ -262,6 +271,35 @@ class Langchain4jChatModelConfigurerImpl implements Langchain4jChatModelConfigur
.build();
}
+ @Override
+ public ChatModel configureChatModel(OllamaChatModelConfig chatModelConfig) {
+ var builder = OllamaChatModel.builder()
+ .baseUrl(chatModelConfig.providerConfig().baseUrl())
+ .modelName(chatModelConfig.modelId())
+ .temperature(chatModelConfig.temperature())
+ .topP(chatModelConfig.topP())
+ .topK(chatModelConfig.topK())
+ .numCtx(chatModelConfig.contextLength())
+ .numPredict(chatModelConfig.maxOutputTokens())
+ .timeout(toDuration(chatModelConfig.timeoutSeconds()))
+ .maxRetries(chatModelConfig.maxRetries());
+
+ var auth = chatModelConfig.providerConfig().auth();
+ if (auth instanceof OllamaProviderConfig.OllamaAuth.Basic basicAuth) {
+ String credentials = basicAuth.username() + ":" + basicAuth.password();
+ String encodedCredentials = Base64.getEncoder().encodeToString(credentials.getBytes(StandardCharsets.UTF_8));
+ builder.customHeaders(singletonMap(HttpHeaders.AUTHORIZATION, "Basic " + encodedCredentials));
+ } else if (auth instanceof OllamaProviderConfig.OllamaAuth.Token tokenAuth) {
+ builder.customHeaders(singletonMap(HttpHeaders.AUTHORIZATION, "Bearer " + tokenAuth.token()));
+ } else if (auth instanceof OllamaProviderConfig.OllamaAuth.None) {
+ // do nothing
+ } else {
+ throw new UnsupportedOperationException("Unknown authentication type: " + auth.getClass().getSimpleName());
+ }
+
+ return builder.build();
+ }
+
private static Duration toDuration(Integer timeoutSeconds) {
return timeoutSeconds != null ? Duration.ofSeconds(timeoutSeconds) : null;
}
diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
index 3dedea999c..84025b3cee 100644
--- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
+++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java
@@ -442,13 +442,13 @@ public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService
boolean check(long threshold, long warnThreshold, long value);
}
- private void checkStartOfNextCycle() {
+ public void checkStartOfNextCycle() {
updateLock.lock();
try {
long now = System.currentTimeMillis();
myUsageStates.values().forEach(state -> {
if ((state.getNextCycleTs() < now) && (now - state.getNextCycleTs() < TimeUnit.HOURS.toMillis(1))) {
- state.setCycles(state.getNextCycleTs(), SchedulerUtils.getStartOfNextNextMonth());
+ state.setCycles(state.getNextCycleTs(), SchedulerUtils.getStartOfNextMonth());
if (log.isTraceEnabled()) {
log.trace("[{}][{}] Updating state cycles (currentCycleTs={},nextCycleTs={})", state.getTenantId(), state.getEntityId(), state.getCurrentCycleTs(), state.getNextCycleTs());
}
diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
new file mode 100644
index 0000000000..fa10a49503
--- /dev/null
+++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
@@ -0,0 +1,252 @@
+/**
+ * Copyright © 2016-2025 The Thingsboard Authors
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.thingsboard.server.service.cf;
+
+import com.google.common.util.concurrent.Futures;
+import com.google.common.util.concurrent.ListenableFuture;
+import com.google.common.util.concurrent.ListeningExecutorService;
+import com.google.common.util.concurrent.MoreExecutors;
+import jakarta.annotation.PostConstruct;
+import jakarta.annotation.PreDestroy;
+import lombok.Data;
+import lombok.extern.slf4j.Slf4j;
+import org.thingsboard.common.util.ThingsBoardExecutors;
+import org.thingsboard.server.common.data.cf.configuration.Argument;
+import org.thingsboard.server.common.data.cf.configuration.ArgumentType;
+import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration;
+import org.thingsboard.server.common.data.id.EntityId;
+import org.thingsboard.server.common.data.id.TenantId;
+import org.thingsboard.server.common.data.kv.Aggregation;
+import org.thingsboard.server.common.data.kv.AttributeKvEntry;
+import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
+import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
+import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
+import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
+import org.thingsboard.server.common.data.kv.TsKvEntry;
+import org.thingsboard.server.common.data.relation.RelationTypeGroup;
+import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
+import org.thingsboard.server.dao.attributes.AttributesService;
+import org.thingsboard.server.dao.relation.RelationService;
+import org.thingsboard.server.dao.timeseries.TimeseriesService;
+import org.thingsboard.server.dao.usagerecord.ApiLimitService;
+import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
+import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
+import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState;
+import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.Set;
+import java.util.concurrent.ExecutionException;
+import java.util.stream.Collectors;
+
+import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY;
+import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY;
+import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry;
+import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createStateByType;
+import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument;
+
+@Data
+@Slf4j
+public abstract class AbstractCalculatedFieldProcessingService {
+
+ protected final AttributesService attributesService;
+ protected final TimeseriesService timeseriesService;
+ protected final ApiLimitService apiLimitService;
+ protected final RelationService relationService;
+
+ protected ListeningExecutorService calculatedFieldCallbackExecutor;
+
+ @PostConstruct
+ public void init() {
+ calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
+ Math.max(4, Runtime.getRuntime().availableProcessors()), getExecutorNamePrefix()));
+ }
+
+ @PreDestroy
+ public void stop() {
+ if (calculatedFieldCallbackExecutor != null) {
+ calculatedFieldCallbackExecutor.shutdownNow();
+ }
+ }
+
+ protected abstract String getExecutorNamePrefix();
+
+ public ListenableFuture fetchStateFromDb(CalculatedFieldCtx ctx, EntityId entityId) {
+ Map> argFutures = switch (ctx.getCalculatedField().getType()) {
+ case GEOFENCING -> fetchGeofencingCalculatedFieldArguments(ctx, entityId, false);
+ case SIMPLE, SCRIPT -> {
+ Map> futures = new HashMap<>();
+ for (var entry : ctx.getArguments().entrySet()) {
+ var argEntityId = resolveEntityId(entityId, entry.getValue());
+ var argValueFuture = fetchArgumentValue(ctx.getTenantId(), argEntityId, entry.getValue(), System.currentTimeMillis());
+ futures.put(entry.getKey(), argValueFuture);
+ }
+ yield futures;
+ }
+ };
+ return Futures.whenAllComplete(argFutures.values()).call(() -> {
+ var result = createStateByType(ctx);
+ result.updateState(ctx, resolveArgumentFutures(argFutures));
+ // TODO: move to state.init() method after merge with alarm rules 2.0
+ if (ctx.hasRelationQueryDynamicArguments() && result instanceof GeofencingCalculatedFieldState geofencingCalculatedFieldState) {
+ geofencingCalculatedFieldState.setLastDynamicArgumentsRefreshTs(System.currentTimeMillis());
+ }
+ return result;
+ }, MoreExecutors.directExecutor());
+ }
+
+ protected EntityId resolveEntityId(EntityId entityId, Argument argument) {
+ return argument.getRefEntityId() != null ? argument.getRefEntityId() : entityId;
+ }
+
+ protected Map resolveArgumentFutures(Map> argFutures) {
+ return argFutures.entrySet().stream()
+ .collect(Collectors.toMap(
+ Map.Entry::getKey, // Keep the key as is
+ entry -> {
+ try {
+ return entry.getValue().get();
+ } catch (ExecutionException e) {
+ Throwable cause = e.getCause();
+ throw new RuntimeException("Failed to fetch " + entry.getKey() + ": " + cause.getMessage(), cause);
+ } catch (InterruptedException e) {
+ throw new RuntimeException("Failed to fetch" + entry.getKey(), e);
+ }
+ }
+ ));
+ }
+
+ protected Map> fetchGeofencingCalculatedFieldArguments(CalculatedFieldCtx ctx, EntityId entityId, boolean dynamicArgumentsOnly) {
+ Map> argFutures = new HashMap<>();
+ Set> entries = ctx.getArguments().entrySet();
+ if (dynamicArgumentsOnly) {
+ entries = entries.stream()
+ .filter(entry -> entry.getValue().hasDynamicSource())
+ .collect(Collectors.toSet());
+ }
+ for (var entry : entries) {
+ switch (entry.getKey()) {
+ case ENTITY_ID_LATITUDE_ARGUMENT_KEY, ENTITY_ID_LONGITUDE_ARGUMENT_KEY ->
+ argFutures.put(entry.getKey(), fetchArgumentValue(ctx.getTenantId(), entityId, entry.getValue(), System.currentTimeMillis()));
+ default -> {
+ var resolvedEntityIdsFuture = resolveGeofencingEntityIds(ctx.getTenantId(), entityId, entry);
+ argFutures.put(entry.getKey(), Futures.transformAsync(resolvedEntityIdsFuture, resolvedEntityIds ->
+ fetchGeofencingKvEntry(ctx.getTenantId(), resolvedEntityIds, entry.getValue()), MoreExecutors.directExecutor()));
+ }
+ }
+ }
+ return argFutures;
+ }
+
+ private ListenableFuture> resolveGeofencingEntityIds(TenantId tenantId, EntityId entityId, Map.Entry entry) {
+ Argument value = entry.getValue();
+ if (value.getRefEntityId() != null) {
+ return Futures.immediateFuture(List.of(value.getRefEntityId()));
+ }
+ if (!value.hasDynamicSource()) {
+ return Futures.immediateFuture(List.of(entityId));
+ }
+ var refDynamicSourceConfiguration = value.getRefDynamicSourceConfiguration();
+ return switch (refDynamicSourceConfiguration.getType()) {
+ case RELATION_PATH_QUERY -> {
+ var configuration = (RelationPathQueryDynamicSourceConfiguration) refDynamicSourceConfiguration;
+ yield Futures.transform(relationService.findByRelationPathQueryAsync(tenantId, configuration.toRelationPathQuery(entityId)),
+ configuration::resolveEntityIds, calculatedFieldCallbackExecutor);
+ }
+ };
+ }
+
+ private ListenableFuture fetchGeofencingKvEntry(TenantId tenantId, List geofencingEntities, Argument argument) {
+ if (argument.getRefEntityKey().getType() != ArgumentType.ATTRIBUTE) {
+ throw new IllegalStateException("Unsupported argument key type: " + argument.getRefEntityKey().getType());
+ }
+ List>> kvFutures = geofencingEntities.stream()
+ .map(entityId -> {
+ var attributesFuture = attributesService.find(
+ tenantId,
+ entityId,
+ argument.getRefEntityKey().getScope(),
+ argument.getRefEntityKey().getKey()
+ );
+ return Futures.transform(attributesFuture, resultOpt ->
+ Map.entry(entityId, resultOpt.orElseGet(() ->
+ new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))),
+ calculatedFieldCallbackExecutor
+ );
+ }).collect(Collectors.toList());
+
+ ListenableFuture>> allFutures = Futures.allAsList(kvFutures);
+
+ return Futures.transform(allFutures, entries -> ArgumentEntry.createGeofencingValueArgument(entries.stream()
+ .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue))), MoreExecutors.directExecutor());
+ }
+
+ protected ListenableFuture fetchArgumentValue(TenantId tenantId, EntityId entityId, Argument argument, long startTs) {
+ return switch (argument.getRefEntityKey().getType()) {
+ case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument, startTs);
+ case ATTRIBUTE -> fetchAttribute(tenantId, entityId, argument, startTs);
+ case TS_LATEST -> fetchTsLatest(tenantId, entityId, argument, startTs);
+ };
+ }
+
+ private ListenableFuture fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument, long queryEndTs) {
+ long argTimeWindow = argument.getTimeWindow() == 0 ? queryEndTs : argument.getTimeWindow();
+ long startInterval = queryEndTs - argTimeWindow;
+ ReadTsKvQuery query = buildTsRollingQuery(tenantId, argument, startInterval, queryEndTs);
+
+ log.trace("[{}][{}] Fetching timeseries for query {}", tenantId, entityId, query);
+ ListenableFuture> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query));
+ return Futures.transform(tsRollingFuture, tsRolling -> {
+ log.debug("[{}][{}] Fetched {} timeseries for query {}", tenantId, entityId, tsRolling == null ? 0 : tsRolling.size(), query);
+ return ArgumentEntry.createTsRollingArgument(tsRolling, query.getLimit(), argTimeWindow);
+ }, calculatedFieldCallbackExecutor);
+ }
+
+ private ListenableFuture fetchAttribute(TenantId tenantId, EntityId entityId, Argument argument, long defaultLastUpdateTs) {
+ log.trace("[{}][{}] Fetching attribute for key {}", tenantId, entityId, argument.getRefEntityKey());
+ var attributeOptFuture = attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey());
+
+ return Futures.transform(attributeOptFuture, attrOpt -> {
+ log.debug("[{}][{}] Fetched attribute for key {}: {}", tenantId, entityId, argument.getRefEntityKey(), attrOpt);
+ AttributeKvEntry attributeKvEntry = attrOpt.orElseGet(() -> new BaseAttributeKvEntry(createDefaultKvEntry(argument), defaultLastUpdateTs, 0L));
+ return transformSingleValueArgument(Optional.of(attributeKvEntry));
+ }, calculatedFieldCallbackExecutor);
+ }
+
+ protected ListenableFuture fetchTsLatest(TenantId tenantId, EntityId entityId, Argument argument, long startTs) {
+ String timeseriesKey = argument.getRefEntityKey().getKey();
+ log.trace("[{}][{}] Fetching latest timeseries {}", tenantId, entityId, timeseriesKey);
+ return transformSingleValueArgument(
+ Futures.transform(
+ timeseriesService.findLatest(tenantId, entityId, timeseriesKey),
+ result -> {
+ log.debug("[{}][{}] Fetched latest timeseries {}: {}", tenantId, entityId, timeseriesKey, result);
+ return result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L)));
+ }, calculatedFieldCallbackExecutor));
+ }
+
+ private ReadTsKvQuery buildTsRollingQuery(TenantId tenantId, Argument argument, long startTs, long endTs) {
+ long maxDataPoints = apiLimitService.getLimit(
+ tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg);
+ int argumentLimit = argument.getLimit();
+ int limit = argumentLimit == 0 || argumentLimit > maxDataPoints ? (int) maxDataPoints : argumentLimit;
+ return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE);
+ }
+
+}
diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
index 847caccaff..86ed174485 100644
--- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
+++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
@@ -34,6 +34,8 @@ public interface CalculatedFieldProcessingService {
ListenableFuture fetchStateFromDb(CalculatedFieldCtx ctx, EntityId entityId);
+ Map fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId);
+
Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments);
void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult calculationResult, List cfIds, TbCallback callback);
diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java
index 49acf6917c..c779c27419 100644
--- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java
+++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java
@@ -34,4 +34,8 @@ public final class CalculatedFieldResult {
(result.isTextual() && result.asText().isEmpty());
}
+ public String toStringOrElseNull() {
+ return result == null ? null : result.toString();
+ }
+
}
diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
index f510857d68..dfe30a0e55 100644
--- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
+++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
@@ -30,6 +30,7 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.page.PageDataIterable;
import org.thingsboard.server.dao.cf.CalculatedFieldService;
+import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.usagerecord.ApiLimitService;
import org.thingsboard.server.queue.util.AfterStartUp;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
@@ -52,6 +53,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
private final CalculatedFieldService calculatedFieldService;
private final TbelInvokeService tbelInvokeService;
private final ApiLimitService apiLimitService;
+ private final RelationService relationService;
private final ConcurrentMap calculatedFields = new ConcurrentHashMap<>();
private final ConcurrentMap> entityIdCalculatedFields = new ConcurrentHashMap<>();
@@ -111,7 +113,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
if (ctx == null) {
CalculatedField calculatedField = getCalculatedField(calculatedFieldId);
if (calculatedField != null) {
- ctx = new CalculatedFieldCtx(calculatedField, tbelInvokeService, apiLimitService);
+ ctx = new CalculatedFieldCtx(calculatedField, tbelInvokeService, apiLimitService, relationService);
calculatedFieldsCtx.put(calculatedFieldId, ctx);
log.debug("[{}] Put calculated field ctx into cache: {}", calculatedFieldId, ctx);
}
diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
index f2a6916751..17dca5dd64 100644
--- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
+++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
@@ -15,46 +15,28 @@
*/
package org.thingsboard.server.service.cf;
-import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
-import com.google.common.util.concurrent.ListeningExecutorService;
-import com.google.common.util.concurrent.MoreExecutors;
-import jakarta.annotation.PostConstruct;
-import jakarta.annotation.PreDestroy;
-import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
-import org.apache.commons.lang3.math.NumberUtils;
import org.springframework.stereotype.Service;
-import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg;
import org.thingsboard.server.actors.calculatedField.MultipleTbCallback;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType;
-import org.thingsboard.server.common.data.StringUtils;
+import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.OutputType;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
-import org.thingsboard.server.common.data.kv.Aggregation;
-import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry;
-import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
-import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
-import org.thingsboard.server.common.data.kv.BooleanDataEntry;
-import org.thingsboard.server.common.data.kv.DoubleDataEntry;
-import org.thingsboard.server.common.data.kv.KvEntry;
-import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
-import org.thingsboard.server.common.data.kv.StringDataEntry;
-import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.msg.TbMsgType;
-import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
import org.thingsboard.server.dao.attributes.AttributesService;
+import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.usagerecord.ApiLimitService;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldLinkedTelemetryMsgProto;
@@ -70,20 +52,12 @@ import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState;
-import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState;
-import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState;
-import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry;
-import org.thingsboard.server.service.cf.ctx.state.TsRollingArgumentEntry;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
-import java.util.Map.Entry;
-import java.util.Optional;
import java.util.UUID;
-import java.util.concurrent.ExecutionException;
-import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.SCOPE;
import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto;
@@ -91,76 +65,53 @@ import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto;
@TbRuleEngineComponent
@Service
@Slf4j
-@RequiredArgsConstructor
-public class DefaultCalculatedFieldProcessingService implements CalculatedFieldProcessingService {
+public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedFieldProcessingService implements CalculatedFieldProcessingService {
- private final AttributesService attributesService;
- private final TimeseriesService timeseriesService;
private final TbClusterService clusterService;
- private final ApiLimitService apiLimitService;
private final PartitionService partitionService;
- private ListeningExecutorService calculatedFieldCallbackExecutor;
-
- @PostConstruct
- public void init() {
- calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool(
- Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback"));
+ public DefaultCalculatedFieldProcessingService(AttributesService attributesService,
+ TimeseriesService timeseriesService,
+ ApiLimitService apiLimitService,
+ RelationService relationService,
+ TbClusterService clusterService,
+ PartitionService partitionService) {
+ super(attributesService, timeseriesService, apiLimitService, relationService);
+ this.clusterService = clusterService;
+ this.partitionService = partitionService;
}
- @PreDestroy
- public void stop() {
- if (calculatedFieldCallbackExecutor != null) {
- calculatedFieldCallbackExecutor.shutdownNow();
- }
+ @Override
+ protected String getExecutorNamePrefix() {
+ return "calculated-field-callback";
}
@Override
public ListenableFuture fetchStateFromDb(CalculatedFieldCtx ctx, EntityId entityId) {
- Map> argFutures = new HashMap<>();
- for (var entry : ctx.getArguments().entrySet()) {
- var argEntityId = entry.getValue().getRefEntityId() != null ? entry.getValue().getRefEntityId() : entityId;
- var argValueFuture = fetchKvEntry(ctx.getTenantId(), argEntityId, entry.getValue());
- argFutures.put(entry.getKey(), argValueFuture);
+ return super.fetchStateFromDb(ctx, entityId);
+ }
+
+ @Override
+ public Map fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId) {
+ // only geofencing calculated fields supports dynamic arguments scheduled updates
+ if (!ctx.getCalculatedField().getType().equals(CalculatedFieldType.GEOFENCING)) {
+ return Map.of();
}
- return Futures.whenAllComplete(argFutures.values()).call(() -> {
- var result = createStateByType(ctx);
- result.updateState(ctx, argFutures.entrySet().stream()
- .collect(Collectors.toMap(
- Entry::getKey, // Keep the key as is
- entry -> {
- try {
- // Resolve the future to get the value
- return entry.getValue().get();
- } catch (ExecutionException | InterruptedException e) {
- throw new RuntimeException("Error getting future result for key: " + entry.getKey(), e);
- }
- }
- )));
- return result;
- }, calculatedFieldCallbackExecutor);
+ return resolveArgumentFutures(fetchGeofencingCalculatedFieldArguments(ctx, entityId, true));
}
@Override
public Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments) {
Map> argFutures = new HashMap<>();
for (var entry : arguments.entrySet()) {
- var argEntityId = entry.getValue().getRefEntityId() != null ? entry.getValue().getRefEntityId() : entityId;
- var argValueFuture = fetchKvEntry(tenantId, argEntityId, entry.getValue());
+ if (entry.getValue().hasDynamicSource()) {
+ continue;
+ }
+ var argEntityId = resolveEntityId(entityId, entry.getValue());
+ var argValueFuture = fetchArgumentValue(tenantId, argEntityId, entry.getValue(), System.currentTimeMillis());
argFutures.put(entry.getKey(), argValueFuture);
}
- return argFutures.entrySet().stream()
- .collect(Collectors.toMap(
- Entry::getKey, // Keep the key as is
- entry -> {
- try {
- // Resolve the future to get the value
- return entry.getValue().get();
- } catch (ExecutionException | InterruptedException e) {
- throw new RuntimeException("Error getting future result for key: " + entry.getKey(), e);
- }
- }
- ));
+ return resolveArgumentFutures(argFutures);
}
@Override
@@ -169,7 +120,7 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
OutputType type = calculatedFieldResult.getType();
TbMsgType msgType = OutputType.ATTRIBUTES.equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST;
TbMsgMetaData md = OutputType.ATTRIBUTES.equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY;
- TbMsg msg = TbMsg.newMsg().type(msgType).originator(entityId).previousCalculatedFieldIds(cfIds).metaData(md).data(calculatedFieldResult.getResult().toString()).build();
+ TbMsg msg = TbMsg.newMsg().type(msgType).originator(entityId).previousCalculatedFieldIds(cfIds).metaData(md).data(calculatedFieldResult.toStringOrElseNull()).build();
clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() {
@Override
public void onSuccess(TbQueueMsgMetadata metadata) {
@@ -241,69 +192,6 @@ public class DefaultCalculatedFieldProcessingService implements CalculatedFieldP
return builder.build();
}
- private ListenableFuture fetchKvEntry(TenantId tenantId, EntityId entityId, Argument argument) {
- return switch (argument.getRefEntityKey().getType()) {
- case TS_ROLLING -> fetchTsRolling(tenantId, entityId, argument);
- case ATTRIBUTE -> transformSingleValueArgument(
- Futures.transform(
- attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()),
- result -> result.or(() -> Optional.of(new BaseAttributeKvEntry(createDefaultKvEntry(argument), System.currentTimeMillis(), 0L))),
- calculatedFieldCallbackExecutor)
- );
- case TS_LATEST -> transformSingleValueArgument(
- Futures.transform(
- timeseriesService.findLatest(tenantId, entityId, argument.getRefEntityKey().getKey()),
- result -> result.or(() -> Optional.of(new BasicTsKvEntry(System.currentTimeMillis(), createDefaultKvEntry(argument), 0L))),
- calculatedFieldCallbackExecutor));
- };
- }
-
- private ListenableFuture transformSingleValueArgument(ListenableFuture> kvEntryFuture) {
- return Futures.transform(kvEntryFuture, kvEntry -> {
- if (kvEntry.isPresent() && kvEntry.get().getValue() != null) {
- return ArgumentEntry.createSingleValueArgument(kvEntry.get());
- } else {
- return new SingleValueArgumentEntry();
- }
- }, calculatedFieldCallbackExecutor);
- }
-
- private ListenableFuture fetchTsRolling(TenantId tenantId, EntityId entityId, Argument argument) {
- long currentTime = System.currentTimeMillis();
- long timeWindow = argument.getTimeWindow() == 0 ? System.currentTimeMillis() : argument.getTimeWindow();
- long startTs = currentTime - timeWindow;
- long maxDataPoints = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg);
- int argumentLimit = argument.getLimit();
- int limit = argumentLimit == 0 || argumentLimit > maxDataPoints ? (int) maxDataPoints : argument.getLimit();
-
- ReadTsKvQuery query = new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, currentTime, 0, limit, Aggregation.NONE);
- ListenableFuture> tsRollingFuture = timeseriesService.findAll(tenantId, entityId, List.of(query));
-
- return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? new TsRollingArgumentEntry(limit, timeWindow) : ArgumentEntry.createTsRollingArgument(tsRolling, limit, timeWindow), calculatedFieldCallbackExecutor);
- }
-
- private KvEntry createDefaultKvEntry(Argument argument) {
- String key = argument.getRefEntityKey().getKey();
- String defaultValue = argument.getDefaultValue();
- if (StringUtils.isBlank(defaultValue)) {
- return new StringDataEntry(key, null);
- }
- if (NumberUtils.isParsable(defaultValue)) {
- return new DoubleDataEntry(key, Double.parseDouble(defaultValue));
- }
- if ("true".equalsIgnoreCase(defaultValue) || "false".equalsIgnoreCase(defaultValue)) {
- return new BooleanDataEntry(key, Boolean.parseBoolean(defaultValue));
- }
- return new StringDataEntry(key, defaultValue);
- }
-
- private CalculatedFieldState createStateByType(CalculatedFieldCtx ctx) {
- return switch (ctx.getCfType()) {
- case SIMPLE -> new SimpleCalculatedFieldState(ctx.getArgNames());
- case SCRIPT -> new ScriptCalculatedFieldState(ctx.getArgNames());
- };
- }
-
private static class TbCallbackWrapper implements TbQueueCallback {
private final TbCallback callback;
diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java
index 83e10b8194..2d43883131 100644
--- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java
+++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java
@@ -19,10 +19,13 @@ import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import org.thingsboard.script.api.tbel.TbelCfArg;
+import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
+import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry;
import java.util.List;
+import java.util.Map;
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
@@ -31,7 +34,8 @@ import java.util.List;
)
@JsonSubTypes({
@JsonSubTypes.Type(value = SingleValueArgumentEntry.class, name = "SINGLE_VALUE"),
- @JsonSubTypes.Type(value = TsRollingArgumentEntry.class, name = "TS_ROLLING")
+ @JsonSubTypes.Type(value = TsRollingArgumentEntry.class, name = "TS_ROLLING"),
+ @JsonSubTypes.Type(value = GeofencingArgumentEntry.class, name = "GEOFENCING")
})
public interface ArgumentEntry {
@@ -58,4 +62,8 @@ public interface ArgumentEntry {
return new TsRollingArgumentEntry(kvEntries, limit, timeWindow);
}
+ static ArgumentEntry createGeofencingValueArgument(Map entityIdkvEntryMap) {
+ return new GeofencingArgumentEntry(entityIdkvEntryMap);
+ }
+
}
diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java
index 68f973c7c1..876bfa2a3f 100644
--- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java
+++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntryType.java
@@ -16,5 +16,5 @@
package org.thingsboard.server.service.cf.ctx.state;
public enum ArgumentEntryType {
- SINGLE_VALUE, TS_ROLLING
+ SINGLE_VALUE, TS_ROLLING, GEOFENCING
}
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 e21d56b6d2..6d877331bd 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
@@ -25,8 +25,6 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
-import static org.thingsboard.server.utils.CalculatedFieldUtils.toSingleValueArgumentProto;
-
@Data
@AllArgsConstructor
public abstract class BaseCalculatedFieldState implements CalculatedFieldState {
@@ -64,7 +62,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState {
boolean entryUpdated;
if (existingEntry == null || newEntry.isForceResetPrevious()) {
- validateNewEntry(newEntry);
+ validateNewEntry(key, newEntry);
arguments.put(key, newEntry);
entryUpdated = true;
} else {
@@ -95,19 +93,7 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState {
}
}
- @Override
- public void checkArgumentSize(String name, ArgumentEntry entry, CalculatedFieldCtx ctx) {
- if (entry instanceof TsRollingArgumentEntry) {
- return;
- }
- if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) {
- if (ctx.getMaxSingleValueArgumentSize() > 0 && toSingleValueArgumentProto(name, singleValueArgumentEntry).getSerializedSize() > ctx.getMaxSingleValueArgumentSize()) {
- throw new IllegalArgumentException("Single value size exceeds the maximum allowed limit. The argument will not be used for calculation.");
- }
- }
- }
-
- protected abstract void validateNewEntry(ArgumentEntry newEntry);
+ protected void validateNewEntry(String key, ArgumentEntry newEntry) {}
private void updateLastUpdateTimestamp(ArgumentEntry entry) {
long newTs = this.latestTimestamp;
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 2e3321eece..13012a028a 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
@@ -25,10 +25,13 @@ import org.thingsboard.server.common.data.cf.CalculatedField;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Argument;
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.ArgumentsBasedCalculatedFieldConfiguration;
+import org.thingsboard.server.common.data.cf.configuration.ExpressionBasedCalculatedFieldConfiguration;
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.ScheduledUpdateSupportedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration;
+import org.thingsboard.server.common.data.cf.configuration.geofencing.GeofencingCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
@@ -36,14 +39,17 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.util.ProtoUtils;
+import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.usagerecord.ApiLimitService;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
+import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.TimeUnit;
import static org.thingsboard.common.util.ExpressionFunctionsUtil.userDefinedFunctions;
@@ -64,6 +70,7 @@ public class CalculatedFieldCtx {
private String expression;
private boolean useLatestTs;
private TbelInvokeService tbelInvokeService;
+ private RelationService relationService;
private CalculatedFieldScriptEngine calculatedFieldScriptEngine;
private ThreadLocal customExpression;
@@ -73,31 +80,63 @@ public class CalculatedFieldCtx {
private long maxStateSize;
private long maxSingleValueArgumentSize;
- public CalculatedFieldCtx(CalculatedField calculatedField, TbelInvokeService tbelInvokeService, ApiLimitService apiLimitService) {
+ private boolean relationQueryDynamicArguments;
+ private List mainEntityGeofencingArgumentNames;
+ private List linkedEntityGeofencingArgumentNames;
+
+ private long scheduledUpdateIntervalMillis;
+
+ public CalculatedFieldCtx(CalculatedField calculatedField, TbelInvokeService tbelInvokeService, ApiLimitService apiLimitService, RelationService relationService) {
this.calculatedField = calculatedField;
this.cfId = calculatedField.getId();
this.tenantId = calculatedField.getTenantId();
this.entityId = calculatedField.getEntityId();
this.cfType = calculatedField.getType();
- CalculatedFieldConfiguration configuration = calculatedField.getConfiguration();
- this.arguments = configuration.getArguments();
+ this.arguments = new HashMap<>();
this.mainEntityArguments = new HashMap<>();
this.linkedEntityArguments = new HashMap<>();
- for (Map.Entry entry : arguments.entrySet()) {
- var refId = entry.getValue().getRefEntityId();
- var refKey = entry.getValue().getRefEntityKey();
- if (refId == null || refId.equals(calculatedField.getEntityId())) {
- mainEntityArguments.put(refKey, entry.getKey());
- } else {
- linkedEntityArguments.computeIfAbsent(refId, key -> new HashMap<>()).put(refKey, entry.getKey());
+ this.argNames = new ArrayList<>();
+ this.mainEntityGeofencingArgumentNames = new ArrayList<>();
+ this.linkedEntityGeofencingArgumentNames = new ArrayList<>();
+ this.output = calculatedField.getConfiguration().getOutput();
+ if (calculatedField.getConfiguration() instanceof ArgumentsBasedCalculatedFieldConfiguration argBasedConfig) {
+ this.arguments.putAll(argBasedConfig.getArguments());
+ for (Map.Entry entry : arguments.entrySet()) {
+ var refId = entry.getValue().getRefEntityId();
+ var refKey = entry.getValue().getRefEntityKey();
+ if (refId == null && entry.getValue().hasDynamicSource()) {
+ relationQueryDynamicArguments = true;
+ continue;
+ }
+ if (refId == null || refId.equals(calculatedField.getEntityId())) {
+ mainEntityArguments.put(refKey, entry.getKey());
+ } else {
+ linkedEntityArguments.computeIfAbsent(refId, key -> new HashMap<>()).put(refKey, entry.getKey());
+ }
+ }
+ this.argNames.addAll(arguments.keySet());
+ if (argBasedConfig instanceof ExpressionBasedCalculatedFieldConfiguration expressionBasedConfig) {
+ this.expression = expressionBasedConfig.getExpression();
+ this.useLatestTs = CalculatedFieldType.SIMPLE.equals(calculatedField.getType()) && ((SimpleCalculatedFieldConfiguration) argBasedConfig).isUseLatestTs();
}
+ if (calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration geofencingConfig) {
+ geofencingConfig.getZoneGroups().forEach((zoneGroupName, config) -> {
+ if (config.isCfEntitySource(entityId)) {
+ mainEntityGeofencingArgumentNames.add(zoneGroupName);
+ return;
+ }
+ if (config.isLinkedCfEntitySource(entityId)) {
+ linkedEntityGeofencingArgumentNames.add(zoneGroupName);
+ }
+ });
+ }
+ }
+ if (calculatedField.getConfiguration() instanceof ScheduledUpdateSupportedCalculatedFieldConfiguration scheduledConfig) {
+ this.scheduledUpdateIntervalMillis = scheduledConfig.isScheduledUpdateEnabled() ? TimeUnit.SECONDS.toMillis(scheduledConfig.getScheduledUpdateInterval()) : -1L;
}
- this.argNames = new ArrayList<>(arguments.keySet());
- this.output = configuration.getOutput();
- this.expression = configuration.getExpression();
- this.useLatestTs = CalculatedFieldType.SIMPLE.equals(calculatedField.getType()) && ((SimpleCalculatedFieldConfiguration) configuration).isUseLatestTs();
this.tbelInvokeService = tbelInvokeService;
+ this.relationService = relationService;
this.maxDataPointsPerRollingArg = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxDataPointsPerRollingArg);
this.maxStateSize = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024;
@@ -105,25 +144,29 @@ public class CalculatedFieldCtx {
}
public void init() {
- if (CalculatedFieldType.SCRIPT.equals(cfType)) {
- try {
- this.calculatedFieldScriptEngine = initEngine(tenantId, expression, tbelInvokeService);
- initialized = true;
- } catch (Exception e) {
- throw new RuntimeException("Failed to init calculated field ctx. Invalid expression syntax.", e);
+ switch (cfType) {
+ case SCRIPT -> {
+ try {
+ this.calculatedFieldScriptEngine = initEngine(tenantId, expression, tbelInvokeService);
+ initialized = true;
+ } catch (Exception e) {
+ throw new RuntimeException("Failed to init calculated field ctx. Invalid expression syntax.", e);
+ }
}
- } else {
- if (isValidExpression(expression)) {
- this.customExpression = ThreadLocal.withInitial(() ->
- new ExpressionBuilder(expression)
- .functions(userDefinedFunctions)
- .implicitMultiplication(true)
- .variables(this.arguments.keySet())
- .build()
- );
- initialized = true;
- } else {
- throw new RuntimeException("Failed to init calculated field ctx. Invalid expression syntax.");
+ case GEOFENCING -> initialized = true;
+ case SIMPLE -> {
+ if (isValidExpression(expression)) {
+ this.customExpression = ThreadLocal.withInitial(() ->
+ new ExpressionBuilder(expression)
+ .functions(userDefinedFunctions)
+ .implicitMultiplication(true)
+ .variables(this.arguments.keySet())
+ .build()
+ );
+ initialized = true;
+ } else {
+ throw new RuntimeException("Failed to init calculated field ctx. Invalid expression syntax.");
+ }
}
}
}
@@ -293,15 +336,42 @@ public class CalculatedFieldCtx {
}
public boolean hasOtherSignificantChanges(CalculatedFieldCtx other) {
- boolean expressionChanged = !expression.equals(other.expression);
+ boolean expressionChanged = calculatedField.getConfiguration() instanceof ExpressionBasedCalculatedFieldConfiguration && !expression.equals(other.expression);
boolean outputChanged = !output.equals(other.output);
- return expressionChanged || outputChanged;
+ boolean scheduledUpdatesConfigChanged = scheduledUpdateIntervalMillis != other.scheduledUpdateIntervalMillis;
+ return expressionChanged || outputChanged || scheduledUpdatesConfigChanged;
}
public boolean hasStateChanges(CalculatedFieldCtx other) {
boolean typeChanged = !cfType.equals(other.cfType);
boolean argumentsChanged = !arguments.equals(other.arguments);
- return typeChanged || argumentsChanged;
+ boolean geoZoneGroupsConfigChanged = hasGeofencingZoneGroupConfigurationChanges(other);
+ return typeChanged || argumentsChanged || geoZoneGroupsConfigChanged;
+ }
+
+ private boolean hasGeofencingZoneGroupConfigurationChanges(CalculatedFieldCtx other) {
+ if (calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration thisConfig
+ && other.calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration otherConfig) {
+ return !thisConfig.getZoneGroups().equals(otherConfig.getZoneGroups());
+ }
+ return false;
+ }
+
+ public boolean hasRelationQueryDynamicArguments() {
+ return relationQueryDynamicArguments && scheduledUpdateIntervalMillis != -1;
+ }
+
+ public boolean shouldFetchDynamicArgumentsFromDb(CalculatedFieldState state) {
+ if (!hasRelationQueryDynamicArguments()) {
+ return false;
+ }
+ if (!(state instanceof GeofencingCalculatedFieldState geofencingState)) {
+ return false;
+ }
+ if (geofencingState.getLastDynamicArgumentsRefreshTs() == -1L) {
+ return true;
+ }
+ return geofencingState.getLastDynamicArgumentsRefreshTs() < System.currentTimeMillis() - scheduledUpdateIntervalMillis;
}
public String getSizeExceedsLimitMessage() {
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 0de354bbb0..3e0964bfd2 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
@@ -20,12 +20,17 @@ import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo;
import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
+import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.service.cf.CalculatedFieldResult;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
+import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingArgumentEntry;
+import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
import java.util.List;
import java.util.Map;
+import static org.thingsboard.server.utils.CalculatedFieldUtils.toSingleValueArgumentProto;
+
@JsonTypeInfo(
use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.PROPERTY,
@@ -34,6 +39,7 @@ import java.util.Map;
@JsonSubTypes({
@JsonSubTypes.Type(value = SimpleCalculatedFieldState.class, name = "SIMPLE"),
@JsonSubTypes.Type(value = ScriptCalculatedFieldState.class, name = "SCRIPT"),
+ @JsonSubTypes.Type(value = GeofencingCalculatedFieldState.class, name = "GEOFENCING"),
})
public interface CalculatedFieldState {
@@ -48,7 +54,7 @@ public interface CalculatedFieldState {
boolean updateState(CalculatedFieldCtx ctx, Map argumentValues);
- ListenableFuture performCalculation(CalculatedFieldCtx ctx);
+ ListenableFuture performCalculation(EntityId entityId, CalculatedFieldCtx ctx);
@JsonIgnore
boolean isReady();
@@ -62,6 +68,15 @@ public interface CalculatedFieldState {
void checkStateSize(CalculatedFieldEntityCtxId ctxId, long maxStateSize);
- void checkArgumentSize(String name, ArgumentEntry entry, CalculatedFieldCtx ctx);
+ default void checkArgumentSize(String name, ArgumentEntry entry, CalculatedFieldCtx ctx) {
+ if (entry instanceof TsRollingArgumentEntry || entry instanceof GeofencingArgumentEntry) {
+ return;
+ }
+ if (entry instanceof SingleValueArgumentEntry singleValueArgumentEntry) {
+ if (ctx.getMaxSingleValueArgumentSize() > 0 && toSingleValueArgumentProto(name, singleValueArgumentEntry).getSerializedSize() > ctx.getMaxSingleValueArgumentSize()) {
+ throw new IllegalArgumentException("Single value size exceeds the maximum allowed limit. The argument will not be used for calculation.");
+ }
+ }
+ }
}
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 84dce627ae..fe7dfa04d0 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
@@ -20,6 +20,7 @@ import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import lombok.Data;
+import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.script.api.tbel.TbelCfArg;
@@ -27,6 +28,7 @@ import org.thingsboard.script.api.tbel.TbelCfCtx;
import org.thingsboard.script.api.tbel.TbelCfSingleValueArg;
import org.thingsboard.server.common.data.cf.CalculatedFieldType;
import org.thingsboard.server.common.data.cf.configuration.Output;
+import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.service.cf.CalculatedFieldResult;
import java.util.ArrayList;
@@ -37,6 +39,7 @@ import java.util.Map;
@Data
@Slf4j
@NoArgsConstructor
+@EqualsAndHashCode(callSuper = true)
public class ScriptCalculatedFieldState extends BaseCalculatedFieldState {
public ScriptCalculatedFieldState(List requiredArguments) {
@@ -49,11 +52,7 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState {
}
@Override
- protected void validateNewEntry(ArgumentEntry newEntry) {
- }
-
- @Override
- public ListenableFuture performCalculation(CalculatedFieldCtx ctx) {
+ public ListenableFuture performCalculation(EntityId entityId, CalculatedFieldCtx ctx) {
Map arguments = new LinkedHashMap<>();
List