|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 21 KiB After Width: | Height: | Size: 19 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 9.0 KiB After Width: | Height: | Size: 7.4 KiB |
|
Before Width: | Height: | Size: 17 KiB After Width: | Height: | Size: 16 KiB |
|
Before Width: | Height: | Size: 10 KiB After Width: | Height: | Size: 8.9 KiB |
@ -0,0 +1,46 @@ |
|||
/** |
|||
* 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.security.auth; |
|||
|
|||
import jakarta.servlet.FilterChain; |
|||
import jakarta.servlet.http.HttpServletRequest; |
|||
import jakarta.servlet.http.HttpServletResponse; |
|||
import lombok.RequiredArgsConstructor; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.springframework.security.core.AuthenticationException; |
|||
import org.springframework.stereotype.Component; |
|||
import org.springframework.web.filter.OncePerRequestFilter; |
|||
import org.thingsboard.server.exception.ThingsboardErrorResponseHandler; |
|||
|
|||
@Component |
|||
@RequiredArgsConstructor |
|||
@Slf4j |
|||
public class AuthExceptionHandler extends OncePerRequestFilter { |
|||
|
|||
private final ThingsboardErrorResponseHandler errorResponseHandler; |
|||
|
|||
@Override |
|||
protected void doFilterInternal(HttpServletRequest request, HttpServletResponse response, FilterChain filterChain) { |
|||
try { |
|||
filterChain.doFilter(request, response); |
|||
} catch (AuthenticationException e) { |
|||
throw e; |
|||
} catch (Exception e) { |
|||
errorResponseHandler.handle(e, response); |
|||
} |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,120 @@ |
|||
/** |
|||
* 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.script.api.js; |
|||
|
|||
import java.util.regex.Pattern; |
|||
|
|||
public class JsValidator { |
|||
|
|||
static final Pattern ASYNC_PATTERN = Pattern.compile("\\basync\\b"); |
|||
static final Pattern AWAIT_PATTERN = Pattern.compile("\\bawait\\b"); |
|||
static final Pattern PROMISE_PATTERN = Pattern.compile("\\bPromise\\b"); |
|||
static final Pattern SET_TIMEOUT_PATTERN = Pattern.compile("\\bsetTimeout\\b"); |
|||
|
|||
public static String validate(String scriptBody) { |
|||
if (scriptBody == null || scriptBody.trim().isEmpty()) { |
|||
return "Script body is empty"; |
|||
} |
|||
|
|||
//Quick check
|
|||
if (!ASYNC_PATTERN.matcher(scriptBody).find() |
|||
&& !AWAIT_PATTERN.matcher(scriptBody).find() |
|||
&& !PROMISE_PATTERN.matcher(scriptBody).find() |
|||
&& !SET_TIMEOUT_PATTERN.matcher(scriptBody).find()) { |
|||
return null; |
|||
} |
|||
|
|||
//Recheck if quick check failed. Ignoring comments and strings
|
|||
String[] lines = scriptBody.split("\\r?\\n"); |
|||
boolean insideMultilineComment = false; |
|||
|
|||
for (String line : lines) { |
|||
String stripped = line; |
|||
|
|||
// Handle multiline comments
|
|||
if (insideMultilineComment) { |
|||
if (line.contains("*/")) { |
|||
insideMultilineComment = false; |
|||
stripped = line.substring(line.indexOf("*/") + 2); // continue after comment
|
|||
} else { |
|||
continue; // skip line inside multiline comment
|
|||
} |
|||
} |
|||
|
|||
// Check for start of multiline comment
|
|||
if (stripped.contains("/*")) { |
|||
int start = stripped.indexOf("/*"); |
|||
int end = stripped.indexOf("*/", start + 2); |
|||
|
|||
if (end != -1) { |
|||
// Inline multiline comment
|
|||
stripped = stripped.substring(0, start) + stripped.substring(end + 2); |
|||
} else { |
|||
// Starts a block comment, continues on next lines
|
|||
insideMultilineComment = true; |
|||
stripped = stripped.substring(0, start); |
|||
} |
|||
} |
|||
|
|||
stripped = stripInlineComment(stripped); |
|||
stripped = stripStringLiterals(stripped); |
|||
|
|||
if (ASYNC_PATTERN.matcher(stripped).find()) { |
|||
return "Script must not contain 'async' keyword."; |
|||
} |
|||
if (AWAIT_PATTERN.matcher(stripped).find()) { |
|||
return "Script must not contain 'await' keyword."; |
|||
} |
|||
if (PROMISE_PATTERN.matcher(stripped).find()) { |
|||
return "Script must not use 'Promise'."; |
|||
} |
|||
if (SET_TIMEOUT_PATTERN.matcher(stripped).find()) { |
|||
return "Script must not use 'setTimeout' method."; |
|||
} |
|||
} |
|||
return null; |
|||
} |
|||
|
|||
private static String stripInlineComment(String line) { |
|||
int index = line.indexOf("//"); |
|||
return index >= 0 ? line.substring(0, index) : line; |
|||
} |
|||
|
|||
private static String stripStringLiterals(String line) { |
|||
StringBuilder sb = new StringBuilder(); |
|||
boolean inSingleQuote = false; |
|||
boolean inDoubleQuote = false; |
|||
|
|||
for (int i = 0; i < line.length(); i++) { |
|||
char c = line.charAt(i); |
|||
|
|||
if (c == '"' && !inSingleQuote) { |
|||
inDoubleQuote = !inDoubleQuote; |
|||
continue; |
|||
} else if (c == '\'' && !inDoubleQuote) { |
|||
inSingleQuote = !inSingleQuote; |
|||
continue; |
|||
} |
|||
|
|||
if (!inSingleQuote && !inDoubleQuote) { |
|||
sb.append(c); |
|||
} |
|||
} |
|||
|
|||
return sb.toString(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,122 @@ |
|||
/** |
|||
* 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.script.api; |
|||
|
|||
import com.google.common.util.concurrent.FutureCallback; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.mockito.Mockito; |
|||
import org.springframework.test.util.ReflectionTestUtils; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.stats.StatsCounter; |
|||
|
|||
import java.util.UUID; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.TimeoutException; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.junit.jupiter.api.Assertions.assertNull; |
|||
import static org.junit.jupiter.api.Assertions.assertThrows; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.anyString; |
|||
import static org.mockito.Mockito.doCallRealMethod; |
|||
import static org.mockito.Mockito.doReturn; |
|||
import static org.mockito.Mockito.mock; |
|||
import static org.mockito.Mockito.never; |
|||
import static org.mockito.Mockito.verify; |
|||
|
|||
class AbstractScriptInvokeServiceTest { |
|||
|
|||
AbstractScriptInvokeService service; |
|||
final UUID id = UUID.randomUUID(); |
|||
final String scriptBody = "return true;"; |
|||
final TenantId tenantId = TenantId.fromUUID(UUID.fromString("2ed9a658-45a5-4812-b212-9931f5749f30")); |
|||
|
|||
@BeforeEach |
|||
void setUp() { |
|||
service = mock(AbstractScriptInvokeService.class, Mockito.RETURNS_DEEP_STUBS); |
|||
|
|||
// Make sure core checks always pass
|
|||
doReturn(true).when(service).isExecEnabled(any()); |
|||
doReturn(50000L).when(service).getMaxScriptBodySize(); |
|||
|
|||
// Use real implementations
|
|||
doCallRealMethod().when(service).scriptBodySizeExceeded(anyString()); |
|||
doCallRealMethod().when(service).eval(any(), any(), any(), any(String[].class)); |
|||
doCallRealMethod().when(service).error(anyString()); |
|||
doCallRealMethod().when(service).validate(any(), anyString()); |
|||
} |
|||
|
|||
@Test |
|||
void evalWithValidationCallTest() throws ExecutionException, InterruptedException, TimeoutException { |
|||
ReflectionTestUtils.setField(service, "requestsCounter", mock(StatsCounter.class)); |
|||
ReflectionTestUtils.setField(service, "evalCallback", mock(FutureCallback.class)); |
|||
|
|||
doReturn(Futures.immediateFuture(id)).when(service).doEvalScript(any(), any(), anyString(), any(), any(String[].class)); |
|||
|
|||
var future = service.eval(tenantId, ScriptType.RULE_NODE_SCRIPT, scriptBody, "x", "y"); |
|||
|
|||
assertThat(future.get(30, TimeUnit.SECONDS)).isEqualTo(id); |
|||
verify(service).validate(any(), anyString()); |
|||
verify(service).validate(tenantId, scriptBody); |
|||
verify(service, never()).error(anyString()); |
|||
} |
|||
|
|||
@Test |
|||
void evalWithValidationCallErrorTest() throws ExecutionException, InterruptedException, TimeoutException { |
|||
doReturn(false).when(service).isExecEnabled(any()); |
|||
var future = service.eval(tenantId, ScriptType.RULE_NODE_SCRIPT, scriptBody, "x", "y"); |
|||
|
|||
ExecutionException ex = assertThrows(ExecutionException.class, future::get); |
|||
assertThat(ex.getCause().getMessage()).isEqualTo("Script Execution is disabled due to API limits!"); |
|||
assertThat(ex.getCause()).isInstanceOf(RuntimeException.class); |
|||
|
|||
verify(service).validate(any(), anyString()); |
|||
verify(service).validate(tenantId, scriptBody); |
|||
verify(service).error(anyString()); |
|||
} |
|||
|
|||
@Test |
|||
void validateScriptBodyTestExecEnabledTest() { |
|||
assertNull(service.validate(tenantId, scriptBody)); |
|||
verify(service).isExecEnabled(tenantId); |
|||
} |
|||
|
|||
@Test |
|||
void validateScriptBodyTestExecDisabledTest() { |
|||
doReturn(false).when(service).isExecEnabled(tenantId); |
|||
assertThat(service.validate(tenantId, scriptBody)).isEqualTo("Script Execution is disabled due to API limits!"); |
|||
verify(service).isExecEnabled(tenantId); |
|||
} |
|||
|
|||
@Test |
|||
void validateScriptBodySizeOKTest() { |
|||
assertNull(service.validate(tenantId, scriptBody)); |
|||
verify(service).isExecEnabled(tenantId); |
|||
verify(service).scriptBodySizeExceeded(scriptBody); |
|||
} |
|||
|
|||
@Test |
|||
void validateScriptBodySizeExceededTest() { |
|||
doReturn(10L).when(service).getMaxScriptBodySize(); |
|||
assertThat(service.validate(tenantId, scriptBody)).isEqualTo("Script body exceeds maximum allowed size of 10 symbols"); |
|||
verify(service).isExecEnabled(tenantId); |
|||
verify(service).scriptBodySizeExceeded(scriptBody); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,91 @@ |
|||
/** |
|||
* 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.script.api.js; |
|||
|
|||
import com.google.common.util.concurrent.FutureCallback; |
|||
import com.google.common.util.concurrent.Futures; |
|||
import org.thingsboard.server.common.stats.StatsCounter; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.mockito.Mockito; |
|||
import org.springframework.test.util.ReflectionTestUtils; |
|||
import org.thingsboard.script.api.ScriptType; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
|
|||
import java.util.UUID; |
|||
import java.util.concurrent.ExecutionException; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.TimeoutException; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.junit.jupiter.api.Assertions.assertThrows; |
|||
import static org.junit.jupiter.api.Assertions.assertTrue; |
|||
import static org.mockito.ArgumentMatchers.any; |
|||
import static org.mockito.ArgumentMatchers.anyString; |
|||
import static org.mockito.Mockito.doCallRealMethod; |
|||
import static org.mockito.Mockito.doReturn; |
|||
import static org.mockito.Mockito.mock; |
|||
import static org.mockito.Mockito.times; |
|||
import static org.mockito.Mockito.verify; |
|||
|
|||
@Slf4j |
|||
class AbstractJsInvokeServiceTest { |
|||
|
|||
AbstractJsInvokeService service; |
|||
final UUID id = UUID.randomUUID(); |
|||
|
|||
@BeforeEach |
|||
void setUp() { |
|||
service = mock(AbstractJsInvokeService.class, Mockito.RETURNS_DEEP_STUBS); |
|||
|
|||
ReflectionTestUtils.setField(service, "requestsCounter", mock(StatsCounter.class)); |
|||
ReflectionTestUtils.setField(service, "evalCallback", mock(FutureCallback.class)); |
|||
|
|||
// Make sure core checks always pass
|
|||
doReturn(true).when(service).isExecEnabled(any()); |
|||
doReturn(false).when(service).scriptBodySizeExceeded(anyString()); |
|||
doReturn(Futures.immediateFuture(id)).when(service).doEvalScript(any(), any(), anyString(), any(), any(String[].class)); |
|||
|
|||
// Use real implementations
|
|||
doCallRealMethod().when(service).eval(any(), any(), any(), any(String[].class)); |
|||
doCallRealMethod().when(service).error(anyString()); |
|||
doCallRealMethod().when(service).validate(any(), anyString()); |
|||
} |
|||
|
|||
@Test |
|||
void shouldReturnValidationErrorFromJsValidator() throws ExecutionException, InterruptedException { |
|||
String scriptWithAsync = "async function test() {}"; |
|||
|
|||
var future = service.eval(TenantId.SYS_TENANT_ID, ScriptType.RULE_NODE_SCRIPT, scriptWithAsync, "a", "b"); |
|||
ExecutionException ex = assertThrows(ExecutionException.class, future::get); |
|||
assertTrue(ex.getCause().getMessage().contains("Script must not contain 'async' keyword.")); |
|||
assertThat(ex.getCause()).isInstanceOf(RuntimeException.class); |
|||
verify(service).isExecEnabled(any()); |
|||
verify(service).scriptBodySizeExceeded(any()); |
|||
} |
|||
|
|||
@Test |
|||
void shouldPassValidationAndCallSuperEval() throws ExecutionException, InterruptedException, TimeoutException { |
|||
String validScript = "function test() { return 42; }"; |
|||
var result = service.eval(TenantId.SYS_TENANT_ID, ScriptType.RULE_NODE_SCRIPT, validScript, "x", "y"); |
|||
|
|||
assertThat(result.get(30, TimeUnit.SECONDS)).isEqualTo(id); |
|||
verify(service, times(1)).isExecEnabled(any()); |
|||
verify(service, times(1)).scriptBodySizeExceeded(any()); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,88 @@ |
|||
/** |
|||
* 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.script.api.js; |
|||
|
|||
import org.junit.jupiter.params.ParameterizedTest; |
|||
import org.junit.jupiter.params.provider.NullAndEmptySource; |
|||
import org.junit.jupiter.params.provider.ValueSource; |
|||
|
|||
import static org.junit.jupiter.api.Assertions.assertEquals; |
|||
import static org.junit.jupiter.api.Assertions.assertNotNull; |
|||
import static org.junit.jupiter.api.Assertions.assertNull; |
|||
|
|||
class JsValidatorTest { |
|||
|
|||
@ParameterizedTest(name = "should return error for script \"{0}\"") |
|||
@ValueSource(strings = { |
|||
"async function test() {}", |
|||
"const result = await someFunc();", |
|||
"const result =\nawait\tsomeFunc();", |
|||
"setTimeout(1000);", |
|||
"new Promise((resolve) => {});", |
|||
"function test() { return 42; } \n\t await test()", |
|||
""" |
|||
function init() { |
|||
await doSomething(); |
|||
} |
|||
""", |
|||
}) |
|||
void shouldReturnErrorForInvalidScripts(String script) { |
|||
assertNotNull(JsValidator.validate(script)); |
|||
} |
|||
|
|||
@ParameterizedTest(name = "should pass validation for script: \"{0}\"") |
|||
@ValueSource(strings = { |
|||
"function test() { return 42; }", |
|||
"const result = 10 * 2;", |
|||
"// async is a keyword but not used: 'const word = \"async\";'", |
|||
"let note = 'setTimeout tight';", |
|||
|
|||
"const word = \"async\";", |
|||
"const word = \"setTimeout\";", |
|||
"const word = \"Promise\";", |
|||
"const word = \"await\";", |
|||
|
|||
"const word = 'async';", |
|||
"const word = 'setTimeout';", |
|||
"const word = 'Promise';", |
|||
"const word = 'await';", |
|||
|
|||
"//function test() { return 42; }", |
|||
"// const result = 10 * 2;", |
|||
"// async is a keyword but not used: 'const word = \"async\";'", |
|||
"//setTimeout(1);", |
|||
|
|||
"a=b+c; // await for a day", |
|||
"return new // Promise((resolve) => {", |
|||
"hello(); // async is a keyword but not used: 'const word = \"async\";'", |
|||
"setGoal(a); //setTimeout(1);", |
|||
|
|||
" /* new Promise((resolve) => {}); // */ return 'await';", |
|||
" /* async */ function calc() {", |
|||
"/* async function abc() { \n await new Promise ( \t setTimeout () ) \n } \n*/", |
|||
}) |
|||
void shouldReturnNullForValidScripts(String script) { |
|||
assertNull(JsValidator.validate(script)); |
|||
} |
|||
|
|||
@ParameterizedTest(name = "should return 'Script body is empty' for input: \"{0}\"") |
|||
@NullAndEmptySource |
|||
@ValueSource(strings = {" ", "\t", "\n"}) |
|||
void shouldReturnErrorForEmptyOrNullScripts(String script) { |
|||
assertEquals("Script body is empty", JsValidator.validate(script)); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,98 @@ |
|||
/** |
|||
* 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.dao.component; |
|||
|
|||
import com.fasterxml.jackson.databind.JsonNode; |
|||
import org.junit.jupiter.api.BeforeEach; |
|||
import org.junit.jupiter.api.Test; |
|||
import org.mockito.Mockito; |
|||
import org.thingsboard.common.util.JacksonUtil; |
|||
import org.thingsboard.server.common.data.id.TenantId; |
|||
import org.thingsboard.server.common.data.plugin.ComponentClusteringMode; |
|||
import org.thingsboard.server.common.data.plugin.ComponentDescriptor; |
|||
import org.thingsboard.server.common.data.plugin.ComponentScope; |
|||
import org.thingsboard.server.common.data.plugin.ComponentType; |
|||
import org.thingsboard.server.dao.exception.IncorrectParameterException; |
|||
|
|||
import static org.junit.jupiter.api.Assertions.assertFalse; |
|||
import static org.junit.jupiter.api.Assertions.assertThrows; |
|||
import static org.junit.jupiter.api.Assertions.assertTrue; |
|||
|
|||
class BaseComponentDescriptorServiceTest { |
|||
|
|||
private BaseComponentDescriptorService service; |
|||
private ComponentDescriptor componentDescriptor; |
|||
private TenantId tenantId; |
|||
|
|||
@BeforeEach |
|||
void setUp() { |
|||
service = Mockito.spy(BaseComponentDescriptorService.class); |
|||
tenantId = TenantId.SYS_TENANT_ID; |
|||
|
|||
// Create a simple component descriptor
|
|||
componentDescriptor = new ComponentDescriptor(); |
|||
componentDescriptor.setType(ComponentType.ACTION); |
|||
componentDescriptor.setScope(ComponentScope.TENANT); |
|||
componentDescriptor.setClusteringMode(ComponentClusteringMode.ENABLED); |
|||
componentDescriptor.setName("Test Component"); |
|||
componentDescriptor.setClazz("org.thingsboard.test.TestComponent"); |
|||
|
|||
// Create configuration descriptor with schema from JSON string
|
|||
String configDescriptorJson = """ |
|||
{ |
|||
"schema": { |
|||
"type": "object", |
|||
"properties": { |
|||
"testField": { |
|||
"type": "string" |
|||
} |
|||
}, |
|||
"required": ["testField"] |
|||
} |
|||
}"""; |
|||
|
|||
componentDescriptor.setConfigurationDescriptor(JacksonUtil.toJsonNode(configDescriptorJson)); |
|||
} |
|||
|
|||
@Test |
|||
void testValidate() { |
|||
// Create valid configuration from JSON string
|
|||
String validConfigJson = "{\"testField\": \"test value\"}"; |
|||
JsonNode validConfig = JacksonUtil.toJsonNode(validConfigJson); |
|||
|
|||
// Create invalid configuration (missing required field) from JSON string
|
|||
String invalidConfigJson = "{}"; |
|||
JsonNode invalidConfig = JacksonUtil.toJsonNode(invalidConfigJson); |
|||
|
|||
// Test valid configuration
|
|||
boolean validResult = service.validate(tenantId, componentDescriptor, validConfig); |
|||
assertTrue(validResult, "Valid configuration should pass validation"); |
|||
|
|||
// Test invalid configuration
|
|||
boolean invalidResult = service.validate(tenantId, componentDescriptor, invalidConfig); |
|||
assertFalse(invalidResult, "Invalid configuration should fail validation"); |
|||
|
|||
// Test with component descriptor without schema
|
|||
ComponentDescriptor noSchemaDescriptor = new ComponentDescriptor(componentDescriptor); |
|||
noSchemaDescriptor.setConfigurationDescriptor(JacksonUtil.toJsonNode("{}")); |
|||
|
|||
// Should throw exception when schema is missing
|
|||
assertThrows(IncorrectParameterException.class, () -> { |
|||
service.validate(tenantId, noSchemaDescriptor, validConfig); |
|||
}, "Should throw exception when schema is missing"); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,21 @@ |
|||
/** |
|||
* 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.mqtt; |
|||
|
|||
@FunctionalInterface |
|||
public interface ReconnectStrategy { |
|||
long getNextReconnectDelay(); |
|||
} |
|||
@ -0,0 +1,82 @@ |
|||
/** |
|||
* 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.mqtt; |
|||
|
|||
import lombok.Getter; |
|||
import lombok.extern.slf4j.Slf4j; |
|||
|
|||
import java.util.concurrent.ThreadLocalRandom; |
|||
import java.util.concurrent.TimeUnit; |
|||
|
|||
@Getter |
|||
@Slf4j |
|||
public class ReconnectStrategyExponential implements ReconnectStrategy { |
|||
|
|||
public static final int DEFAULT_RECONNECT_INTERVAL_SEC = 10; |
|||
public static final int MAX_RECONNECT_INTERVAL_SEC = 60; |
|||
public static final int EXP_MAX = 8; |
|||
public static final long JITTER_MAX = 1; |
|||
private final long reconnectIntervalMinSeconds; |
|||
private final long reconnectIntervalMaxSeconds; |
|||
private long lastDisconnectNanoTime = 0; //isotonic time
|
|||
private long retryCount = 0; |
|||
|
|||
public ReconnectStrategyExponential(long reconnectIntervalMinSeconds) { |
|||
this.reconnectIntervalMaxSeconds = calculateIntervalMax(reconnectIntervalMinSeconds); |
|||
this.reconnectIntervalMinSeconds = calculateIntervalMin(reconnectIntervalMinSeconds); |
|||
} |
|||
|
|||
long calculateIntervalMax(long reconnectIntervalMinSeconds) { |
|||
return reconnectIntervalMinSeconds > MAX_RECONNECT_INTERVAL_SEC ? reconnectIntervalMinSeconds : MAX_RECONNECT_INTERVAL_SEC; |
|||
} |
|||
|
|||
long calculateIntervalMin(long reconnectIntervalMinSeconds) { |
|||
return Math.min((reconnectIntervalMinSeconds > 0 ? reconnectIntervalMinSeconds : DEFAULT_RECONNECT_INTERVAL_SEC), this.reconnectIntervalMaxSeconds); |
|||
} |
|||
|
|||
@Override |
|||
synchronized public long getNextReconnectDelay() { |
|||
final long currentNanoTime = getNanoTime(); |
|||
final long coolDownSpentNanos = currentNanoTime - lastDisconnectNanoTime; |
|||
lastDisconnectNanoTime = currentNanoTime; |
|||
if (isCooledDown(coolDownSpentNanos)) { |
|||
retryCount = 0; |
|||
return reconnectIntervalMinSeconds; |
|||
} |
|||
return calculateNextReconnectDelay() + calculateJitter(); |
|||
} |
|||
|
|||
long calculateJitter() { |
|||
return ThreadLocalRandom.current().nextInt() >= 0 ? JITTER_MAX : 0; |
|||
} |
|||
|
|||
long calculateNextReconnectDelay() { |
|||
return Math.min(reconnectIntervalMaxSeconds, reconnectIntervalMinSeconds + calculateExp(retryCount++)); |
|||
} |
|||
|
|||
long calculateExp(long e) { |
|||
return 1L << Math.min(e, EXP_MAX); |
|||
} |
|||
|
|||
boolean isCooledDown(long coolDownSpentNanos) { |
|||
return TimeUnit.NANOSECONDS.toSeconds(coolDownSpentNanos) > reconnectIntervalMaxSeconds + reconnectIntervalMinSeconds; |
|||
} |
|||
|
|||
long getNanoTime() { |
|||
return System.nanoTime(); |
|||
} |
|||
|
|||
} |
|||
@ -0,0 +1,95 @@ |
|||
/** |
|||
* 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.mqtt; |
|||
|
|||
import lombok.extern.slf4j.Slf4j; |
|||
import org.junit.jupiter.api.parallel.Execution; |
|||
import org.junit.jupiter.api.parallel.ExecutionMode; |
|||
import org.junit.jupiter.params.ParameterizedTest; |
|||
import org.junit.jupiter.params.provider.ValueSource; |
|||
import org.mockito.Mockito; |
|||
import org.mockito.stubbing.Answer; |
|||
|
|||
import java.util.ArrayList; |
|||
import java.util.Collection; |
|||
import java.util.concurrent.BlockingQueue; |
|||
import java.util.concurrent.LinkedBlockingDeque; |
|||
import java.util.concurrent.TimeUnit; |
|||
import java.util.concurrent.atomic.AtomicLong; |
|||
|
|||
import static org.assertj.core.api.Assertions.assertThat; |
|||
import static org.assertj.core.data.Offset.offset; |
|||
import static org.mockito.ArgumentMatchers.anyLong; |
|||
import static org.mockito.BDDMockito.willAnswer; |
|||
import static org.thingsboard.mqtt.ReconnectStrategyExponential.EXP_MAX; |
|||
import static org.thingsboard.mqtt.ReconnectStrategyExponential.JITTER_MAX; |
|||
|
|||
@Slf4j |
|||
class ReconnectStrategyExponentialTest { |
|||
|
|||
@Execution(ExecutionMode.SAME_THREAD) // just for convenient log reading
|
|||
@ParameterizedTest |
|||
@ValueSource(ints = {1, 0, 60}) |
|||
public void exponentialReconnectDelayTest(final int reconnectIntervalMinSeconds) { |
|||
final ReconnectStrategyExponential strategy = Mockito.spy(new ReconnectStrategyExponential(reconnectIntervalMinSeconds)); |
|||
log.info("=== Reconnect delay test for ReconnectStrategyExponential({}) : calculated min [{}] max [{}] ===", reconnectIntervalMinSeconds, strategy.getReconnectIntervalMinSeconds(), strategy.getReconnectIntervalMaxSeconds()); |
|||
final AtomicLong nanoTime = new AtomicLong(System.nanoTime()); |
|||
willAnswer((x) -> nanoTime.get()).given(strategy).getNanoTime(); |
|||
final LinkedBlockingDeque<Long> jittersCaptured = new LinkedBlockingDeque<>(); |
|||
final LinkedBlockingDeque<Long> expCaptured = new LinkedBlockingDeque<>(); |
|||
|
|||
willAnswer(captureResult(jittersCaptured)).given(strategy).calculateJitter(); |
|||
willAnswer(captureResult(expCaptured)).given(strategy).calculateExp(anyLong()); |
|||
|
|||
for (int phase = 0; phase < 3; phase++) { |
|||
log.info("== Phase {} ==", phase); |
|||
long previousDelay = 0; |
|||
for (int i = 0; i < EXP_MAX + 4; i++) { |
|||
final long nextReconnectDelay = strategy.getNextReconnectDelay(); |
|||
nanoTime.addAndGet(TimeUnit.SECONDS.toNanos(nextReconnectDelay)); |
|||
log.info("Retry [{}] Delay [{}] : min [{}] exp [{}] jitter [{}]", strategy.getRetryCount(), nextReconnectDelay, strategy.getReconnectIntervalMinSeconds(), expCaptured.peekLast(), jittersCaptured.peekLast()); |
|||
assertThat(previousDelay).satisfiesAnyOf( |
|||
v -> assertThat(v).isLessThanOrEqualTo(nextReconnectDelay), |
|||
v -> assertThat(v).isCloseTo(nextReconnectDelay, offset(JITTER_MAX)) // Adjust tolerance as needed
|
|||
); |
|||
previousDelay = nextReconnectDelay; |
|||
} |
|||
log.info("Jitters captured: {}", drainAll(jittersCaptured)); |
|||
log.info("Exponents captured: {}", drainAll(expCaptured)); |
|||
assertThat(previousDelay).isCloseTo(strategy.getReconnectIntervalMaxSeconds(), offset(JITTER_MAX)); |
|||
|
|||
final long coolDownPeriodSec = strategy.getReconnectIntervalMinSeconds() + strategy.getReconnectIntervalMaxSeconds() + 1; |
|||
log.info("Cooling down for [{}] seconds ...", coolDownPeriodSec); |
|||
nanoTime.addAndGet(TimeUnit.SECONDS.toNanos(coolDownPeriodSec)); |
|||
assertThat(strategy.isCooledDown(TimeUnit.SECONDS.toNanos(coolDownPeriodSec))).as("cooled down").isTrue(); |
|||
} |
|||
} |
|||
|
|||
private Answer<Long> captureResult(Collection<Long> collection) { |
|||
return invocation -> { |
|||
long result = (long) invocation.callRealMethod(); |
|||
collection.add(result); |
|||
return result; |
|||
}; |
|||
} |
|||
|
|||
private Collection<Long> drainAll(BlockingQueue<Long> jittersCaptured) { |
|||
Collection<Long> elements = new ArrayList<>(); |
|||
jittersCaptured.drainTo(elements); |
|||
return elements; |
|||
} |
|||
|
|||
} |
|||