88 changed files with 1741 additions and 698 deletions
|
Before Width: | Height: | Size: 56 KiB After Width: | Height: | Size: 56 KiB |
|
Before Width: | Height: | Size: 56 KiB After Width: | Height: | Size: 56 KiB |
@ -0,0 +1,16 @@ |
|||||
|
-- |
||||
|
-- Copyright © 2016-2024 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. |
||||
|
-- |
||||
|
|
||||
@ -0,0 +1,126 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 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.subscription; |
||||
|
|
||||
|
import ch.qos.logback.classic.Logger; |
||||
|
import ch.qos.logback.classic.spi.ILoggingEvent; |
||||
|
import ch.qos.logback.core.read.ListAppender; |
||||
|
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 org.apache.commons.lang3.RandomStringUtils; |
||||
|
import org.junit.jupiter.api.AfterEach; |
||||
|
import org.junit.jupiter.api.BeforeEach; |
||||
|
import org.junit.jupiter.api.Test; |
||||
|
import org.slf4j.LoggerFactory; |
||||
|
import org.springframework.test.util.ReflectionTestUtils; |
||||
|
import org.thingsboard.server.cache.limits.RateLimitService; |
||||
|
import org.thingsboard.server.common.data.id.DeviceId; |
||||
|
import org.thingsboard.server.common.data.id.EntityId; |
||||
|
import org.thingsboard.server.common.data.id.TenantId; |
||||
|
import org.thingsboard.server.common.data.limit.LimitedApi; |
||||
|
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; |
||||
|
import org.thingsboard.server.gen.transport.TransportProtos; |
||||
|
import org.thingsboard.server.queue.discovery.PartitionService; |
||||
|
import org.thingsboard.server.service.ws.WebSocketSessionRef; |
||||
|
|
||||
|
import java.util.ArrayList; |
||||
|
import java.util.HashMap; |
||||
|
import java.util.List; |
||||
|
import java.util.Map; |
||||
|
import java.util.UUID; |
||||
|
import java.util.concurrent.Executors; |
||||
|
|
||||
|
import static org.junit.jupiter.api.Assertions.assertFalse; |
||||
|
import static org.mockito.ArgumentMatchers.any; |
||||
|
import static org.mockito.ArgumentMatchers.eq; |
||||
|
import static org.mockito.ArgumentMatchers.nullable; |
||||
|
import static org.mockito.Mockito.mock; |
||||
|
import static org.mockito.Mockito.when; |
||||
|
|
||||
|
public class DefaultTbLocalSubscriptionServiceTest { |
||||
|
|
||||
|
ListAppender<ILoggingEvent> testLogAppender; |
||||
|
TbLocalSubscriptionService subscriptionService; |
||||
|
|
||||
|
@BeforeEach |
||||
|
public void setUp() throws Exception { |
||||
|
Logger logger = (Logger) LoggerFactory.getLogger(DefaultTbLocalSubscriptionService.class); |
||||
|
testLogAppender = new ListAppender<>(); |
||||
|
testLogAppender.start(); |
||||
|
logger.addAppender(testLogAppender); |
||||
|
|
||||
|
RateLimitService rateLimitService = mock(); |
||||
|
when(rateLimitService.checkRateLimit(eq(LimitedApi.WS_SUBSCRIPTIONS), any(Object.class), nullable(String.class))).thenReturn(true); |
||||
|
PartitionService partitionService = mock(); |
||||
|
when(partitionService.resolve(any(), any(), any())).thenReturn(TopicPartitionInfo.builder().build()); |
||||
|
subscriptionService = new DefaultTbLocalSubscriptionService(mock(), mock(), mock(), partitionService, mock(), mock(), mock(), rateLimitService); |
||||
|
ReflectionTestUtils.setField(subscriptionService, "serviceId", "serviceId"); |
||||
|
} |
||||
|
|
||||
|
@AfterEach |
||||
|
public void tearDown() { |
||||
|
if (testLogAppender != null) { |
||||
|
testLogAppender.stop(); |
||||
|
Logger logger = (Logger) LoggerFactory.getLogger(DefaultTbLocalSubscriptionService.class); |
||||
|
logger.detachAppender(testLogAppender); |
||||
|
} |
||||
|
} |
||||
|
|
||||
|
@Test |
||||
|
public void addSubscriptionConcurrentModificationTest() throws Exception { |
||||
|
ListeningExecutorService executorService = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(10)); |
||||
|
TenantId tenantId = new TenantId(UUID.randomUUID()); |
||||
|
DeviceId deviceId = new DeviceId(UUID.randomUUID()); |
||||
|
WebSocketSessionRef sessionRef = mock(); |
||||
|
ReflectionTestUtils.setField(subscriptionService, "subscriptionUpdateExecutor", executorService); |
||||
|
|
||||
|
List<ListenableFuture<?>> futures = new ArrayList<>(); |
||||
|
|
||||
|
try { |
||||
|
subscriptionService.onCoreStartupMsg(TransportProtos.CoreStartupMsg.newBuilder().addAllPartitions(List.of(0)).getDefaultInstanceForType()); |
||||
|
for (int i = 0; i < 50; i++) { |
||||
|
futures.add(executorService.submit(() -> subscriptionService.addSubscription(createSubscription(tenantId, deviceId), sessionRef))); |
||||
|
} |
||||
|
Futures.allAsList(futures).get(); |
||||
|
} finally { |
||||
|
executorService.shutdownNow(); |
||||
|
} |
||||
|
|
||||
|
List<ILoggingEvent> logs = testLogAppender.list; |
||||
|
boolean exceptionLogged = logs.stream() |
||||
|
.filter(event -> event.getThrowableProxy() != null) |
||||
|
.map(event -> event.getThrowableProxy().getClassName()) |
||||
|
.anyMatch(log -> log.equals("java.util.ConcurrentModificationException")); |
||||
|
|
||||
|
assertFalse(exceptionLogged, "Detected ConcurrentModificationException!"); |
||||
|
} |
||||
|
|
||||
|
private TbSubscription<?> createSubscription(TenantId tenantId, EntityId entityId) { |
||||
|
Map<String, Long> keys = new HashMap<>(); |
||||
|
for (int i = 0; i < 50; i++) { |
||||
|
keys.put(RandomStringUtils.randomAlphanumeric(5), 1L); |
||||
|
} |
||||
|
return TbAttributeSubscription.builder() |
||||
|
.tenantId(tenantId) |
||||
|
.entityId(entityId) |
||||
|
.subscriptionId(1) |
||||
|
.sessionId(RandomStringUtils.randomAlphanumeric(5)) |
||||
|
.keyStates(keys) |
||||
|
.build(); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,40 @@ |
|||||
|
/** |
||||
|
* Copyright © 2016-2024 The Thingsboard Authors |
||||
|
* |
||||
|
* Licensed under the Apache License, Version 2.0 (the "License"); |
||||
|
* you may not use this file except in compliance with the License. |
||||
|
* You may obtain a copy of the License at |
||||
|
* |
||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
* |
||||
|
* Unless required by applicable law or agreed to in writing, software |
||||
|
* distributed under the License is distributed on an "AS IS" BASIS, |
||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
||||
|
* See the License for the specific language governing permissions and |
||||
|
* limitations under the License. |
||||
|
*/ |
||||
|
package org.thingsboard.server.transport.coap.efento.utils; |
||||
|
|
||||
|
public enum PulseCounterType { |
||||
|
|
||||
|
WATER_CNT_ACC("water_cnt_acc_", 100), |
||||
|
PULSE_CNT_ACC("pulse_cnt_acc_", 1000), |
||||
|
ELEC_METER_ACC("elec_meter_acc_", 1000), |
||||
|
PULSE_CNT_ACC_WIDE("pulse_cnt_acc_wide_", 1000000); |
||||
|
|
||||
|
private final String prefix; |
||||
|
private final int majorResolution; |
||||
|
|
||||
|
PulseCounterType(String prefix, int majorResolution) { |
||||
|
this.prefix = prefix; |
||||
|
this.majorResolution = majorResolution; |
||||
|
} |
||||
|
|
||||
|
public String getPrefix() { |
||||
|
return prefix; |
||||
|
} |
||||
|
|
||||
|
public int getMajorResolution() { |
||||
|
return majorResolution; |
||||
|
} |
||||
|
} |
||||
File diff suppressed because one or more lines are too long
@ -0,0 +1,130 @@ |
|||||
|
/* |
||||
|
* Copyright © 2016-2024 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. |
||||
|
*/ |
||||
|
|
||||
|
const fs = require('fs'); |
||||
|
const path = require('path'); |
||||
|
|
||||
|
const materialIconDir = path.join('.', 'src', 'assets', 'metadata'); |
||||
|
const mdiMetadata = path.join('.', 'node_modules', '@mdi', 'svg', 'meta.json'); |
||||
|
|
||||
|
async function init() { |
||||
|
const iconsBundle = JSON.parse(await fs.promises.readFile(path.join(materialIconDir, 'material-icons.json'))); |
||||
|
|
||||
|
await getMaterialIconMetadataAndUpdated(iconsBundle); |
||||
|
await getMDIMetadataAndUpdated(iconsBundle); |
||||
|
|
||||
|
await fs.promises.writeFile(path.join(materialIconDir, 'material-icons.json'), JSON.stringify(iconsBundle), 'utf8') |
||||
|
} |
||||
|
|
||||
|
async function getMaterialIconMetadataAndUpdated(iconsBundle){ |
||||
|
const iconsResponse = await fetch('https://fonts.google.com/metadata/icons?key=material_symbols&incomplete=true'); |
||||
|
const iconsText = await iconsResponse.text(); |
||||
|
const clearText = iconsText.substring(iconsText.indexOf("\n") + 1); |
||||
|
|
||||
|
const icons = JSON.parse(clearText).icons; |
||||
|
|
||||
|
let prevItem; |
||||
|
const filterIcons = icons.filter((item) => { |
||||
|
if (prevItem?.name !== item.name && !item.unsupported_families.includes('Material Icons')) { |
||||
|
prevItem = item; |
||||
|
return true; |
||||
|
} |
||||
|
return false; |
||||
|
}); |
||||
|
|
||||
|
filterIcons.forEach((item, index) => { |
||||
|
const findItem = iconsBundle.find((el) => el.name === item.name); |
||||
|
if (!findItem) { |
||||
|
let prevIndexIcon = 0; |
||||
|
if (index === 0) { |
||||
|
prevIndexIcon = 45; |
||||
|
} else { |
||||
|
let iteration = 0; |
||||
|
while (prevIndexIcon < 45) { |
||||
|
iteration++; |
||||
|
const prevIconName = filterIcons[index - iteration].name; |
||||
|
prevIndexIcon = findPreviousIcon(iconsBundle, prevIconName); |
||||
|
} |
||||
|
} |
||||
|
if (prevIndexIcon >= 0) { |
||||
|
iconsBundle.splice(prevIndexIcon + 1, 0, {name:item.name, tags:item.tags}); |
||||
|
} |
||||
|
console.log('Not found icon:', item.name); |
||||
|
console.count('Not found material icon'); |
||||
|
return; |
||||
|
} |
||||
|
if (JSON.stringify(item.tags) !== JSON.stringify(findItem.tags)) { |
||||
|
findItem.tags = item.tags; |
||||
|
console.log('Difference tags in', item.name); |
||||
|
console.count('Difference tags in material icon'); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
async function getMDIMetadataAndUpdated(iconsBundle){ |
||||
|
const mdiBundle = JSON.parse(await fs.promises.readFile(mdiMetadata)); |
||||
|
|
||||
|
iconsBundle |
||||
|
.filter(item => item.name.startsWith('mdi:')) |
||||
|
.forEach(item => { |
||||
|
const iconName = item.name.substring(item.name.indexOf(":") + 1); |
||||
|
const findItem = mdiBundle.find((el) => el.name === iconName); |
||||
|
if (!findItem) { |
||||
|
console.error('Delete icon:', item.name); |
||||
|
} |
||||
|
}); |
||||
|
|
||||
|
|
||||
|
mdiBundle.forEach((item, index) => { |
||||
|
const iconName = `mdi:${item.name}` |
||||
|
let iconTags = item.tags; |
||||
|
const iconAliases = item.aliases.map(item => item.replaceAll('-', ' ')); |
||||
|
if (!iconTags.length && item.aliases.length) { |
||||
|
iconTags = iconAliases; |
||||
|
} else if (item.aliases.length) { |
||||
|
iconTags = iconTags.concat(iconAliases); |
||||
|
} |
||||
|
iconTags = iconTags.map(item => item.toLowerCase()); |
||||
|
|
||||
|
const findItem = iconsBundle.find((el) => el.name === iconName); |
||||
|
if (!findItem) { |
||||
|
let prevIndexIcon; |
||||
|
if (index === 0) { |
||||
|
prevIndexIcon = iconsBundle.findIndex(item => item.name.startsWith('mdi:')) |
||||
|
} else { |
||||
|
const prevIconName = `mdi:${mdiBundle[index - 1].name}`; |
||||
|
prevIndexIcon = findPreviousIcon(iconsBundle, prevIconName); |
||||
|
} |
||||
|
if (prevIndexIcon >= 0) { |
||||
|
iconsBundle.splice(prevIndexIcon + 1, 0, {name:iconName, tags:iconTags}); |
||||
|
} |
||||
|
console.log('Not found icon:', iconName); |
||||
|
console.count('Not found mdi icon'); |
||||
|
return; |
||||
|
} |
||||
|
if (JSON.stringify(iconTags) !== JSON.stringify(findItem.tags)) { |
||||
|
findItem.tags = iconTags; |
||||
|
console.log('Difference tags in', iconName); |
||||
|
console.count('Difference tags in mdi icon'); |
||||
|
} |
||||
|
}); |
||||
|
} |
||||
|
|
||||
|
function findPreviousIcon(iconsBundle, findName) { |
||||
|
return iconsBundle.findIndex(item => item.name === findName); |
||||
|
} |
||||
|
|
||||
|
init(); |
||||
@ -0,0 +1,32 @@ |
|||||
|
///
|
||||
|
/// Copyright © 2016-2024 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.
|
||||
|
///
|
||||
|
|
||||
|
import { Pipe, PipeTransform } from '@angular/core'; |
||||
|
import { GatewayVersion } from '@home/components/widget/lib/gateway/gateway-widget.models'; |
||||
|
import { |
||||
|
GatewayConnectorVersionMappingUtil |
||||
|
} from '@home/components/widget/lib/gateway/utils/gateway-connector-version-mapping.util'; |
||||
|
|
||||
|
@Pipe({ |
||||
|
name: 'isLatestVersionConfig', |
||||
|
standalone: true, |
||||
|
}) |
||||
|
export class LatestVersionConfigPipe implements PipeTransform { |
||||
|
transform(configVersion: number | string): boolean { |
||||
|
return GatewayConnectorVersionMappingUtil.parseVersion(configVersion) |
||||
|
>= GatewayConnectorVersionMappingUtil.parseVersion(GatewayVersion.Current); |
||||
|
} |
||||
|
} |
||||
@ -0,0 +1,65 @@ |
|||||
|
### Expressions |
||||
|
#### JSON Path: |
||||
|
|
||||
|
The expression field is used to extract data from the MQTT message. There are various available options for different parts of the messages: |
||||
|
|
||||
|
- The JSONPath format can be used to extract data from the message body. |
||||
|
|
||||
|
- The regular expression format can be used to extract data from the topic where the message will arrive. |
||||
|
|
||||
|
- Slices can only be used in the expression fields of bytes converters. |
||||
|
|
||||
|
JSONPath expressions specify the items within a JSON structure (which could be an object, array, or nested combination of both) that you want to access. These expressions can select elements from JSON data on specific criteria. Here's a basic overview of how JSONPath expressions are structured: |
||||
|
|
||||
|
- `$`: The root element of the JSON document; |
||||
|
|
||||
|
- `.`: Child operator used to select child elements. For example, $.store.book ; |
||||
|
|
||||
|
- `[]`: Child operator used to select child elements. $['store']['book'] accesses the book array within a store object; |
||||
|
|
||||
|
##### Examples: |
||||
|
|
||||
|
For example, if we want to extract the device name from the following message, we can use the expression below: |
||||
|
|
||||
|
MQTT message: |
||||
|
|
||||
|
``` |
||||
|
{ |
||||
|
"sensorModelInfo": { |
||||
|
"sensorName": "AM-123", |
||||
|
"sensorType": "myDeviceType" |
||||
|
}, |
||||
|
"data": { |
||||
|
"temp": 12.2, |
||||
|
"hum": 56, |
||||
|
"status": "ok" |
||||
|
} |
||||
|
} |
||||
|
{:copy-code} |
||||
|
``` |
||||
|
|
||||
|
Expression: |
||||
|
|
||||
|
`${sensorModelInfo.sensorName}` |
||||
|
|
||||
|
Converted data: |
||||
|
|
||||
|
`AM-123` |
||||
|
|
||||
|
If we want to extract all data from the message above, we can use the following expression: |
||||
|
|
||||
|
`${data}` |
||||
|
|
||||
|
Converted data: |
||||
|
|
||||
|
`{"temp": 12.2, "hum": 56, "status": "ok"}` |
||||
|
|
||||
|
Or if we want to extract specific data (for example “temperature”), you can use the following expression: |
||||
|
|
||||
|
`${data.temp}` |
||||
|
|
||||
|
And as a converted data we will get: |
||||
|
|
||||
|
`12.2` |
||||
|
|
||||
|
<br/> |
||||
File diff suppressed because one or more lines are too long
Loading…
Reference in new issue