|
|
|
@ -16,7 +16,6 @@ |
|
|
|
package org.thingsboard.server.service.telemetry; |
|
|
|
|
|
|
|
import com.fasterxml.jackson.core.JsonProcessingException; |
|
|
|
import com.fasterxml.jackson.databind.ObjectMapper; |
|
|
|
import com.google.common.base.Function; |
|
|
|
import com.google.common.util.concurrent.FutureCallback; |
|
|
|
import com.google.common.util.concurrent.Futures; |
|
|
|
@ -27,6 +26,7 @@ import org.springframework.beans.factory.annotation.Autowired; |
|
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
|
import org.springframework.stereotype.Service; |
|
|
|
import org.springframework.web.socket.CloseStatus; |
|
|
|
import org.thingsboard.common.util.JacksonUtil; |
|
|
|
import org.thingsboard.common.util.ThingsBoardExecutors; |
|
|
|
import org.thingsboard.common.util.ThingsBoardThreadFactory; |
|
|
|
import org.thingsboard.server.common.data.DataConstants; |
|
|
|
@ -95,6 +95,8 @@ import java.util.concurrent.ExecutorService; |
|
|
|
import java.util.concurrent.Executors; |
|
|
|
import java.util.concurrent.ScheduledExecutorService; |
|
|
|
import java.util.concurrent.TimeUnit; |
|
|
|
import java.util.concurrent.locks.Lock; |
|
|
|
import java.util.concurrent.locks.ReentrantLock; |
|
|
|
import java.util.function.Consumer; |
|
|
|
import java.util.stream.Collectors; |
|
|
|
|
|
|
|
@ -112,7 +114,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
private static final Aggregation DEFAULT_AGGREGATION = Aggregation.NONE; |
|
|
|
private static final int UNKNOWN_SUBSCRIPTION_ID = 0; |
|
|
|
private static final String PROCESSING_MSG = "[{}] Processing: {}"; |
|
|
|
private static final ObjectMapper jsonMapper = new ObjectMapper(); |
|
|
|
private static final String FAILED_TO_FETCH_DATA = "Failed to fetch data!"; |
|
|
|
private static final String FAILED_TO_FETCH_ATTRIBUTES = "Failed to fetch attributes!"; |
|
|
|
private static final String SESSION_META_DATA_NOT_FOUND = "Session meta-data not found!"; |
|
|
|
@ -147,10 +148,10 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
@Value("${server.ws.ping_timeout:30000}") |
|
|
|
private long pingTimeout; |
|
|
|
|
|
|
|
private ConcurrentMap<TenantId, Set<String>> tenantSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
private ConcurrentMap<CustomerId, Set<String>> customerSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
private ConcurrentMap<UserId, Set<String>> regularUserSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
private ConcurrentMap<UserId, Set<String>> publicUserSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<TenantId, Set<String>> tenantSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<CustomerId, Set<String>> customerSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<UserId, Set<String>> regularUserSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<UserId, Set<String>> publicUserSubscriptionsMap = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
private ExecutorService executor; |
|
|
|
private String serviceId; |
|
|
|
@ -204,7 +205,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
} |
|
|
|
|
|
|
|
try { |
|
|
|
TelemetryPluginCmdsWrapper cmdsWrapper = jsonMapper.readValue(msg, TelemetryPluginCmdsWrapper.class); |
|
|
|
TelemetryPluginCmdsWrapper cmdsWrapper = JacksonUtil.OBJECT_MAPPER.readValue(msg, TelemetryPluginCmdsWrapper.class); |
|
|
|
if (cmdsWrapper != null) { |
|
|
|
if (cmdsWrapper.getAttrSubCmds() != null) { |
|
|
|
cmdsWrapper.getAttrSubCmds().forEach(cmd -> { |
|
|
|
@ -450,7 +451,6 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
@Override |
|
|
|
public void onSuccess(List<AttributeKvEntry> data) { |
|
|
|
List<TsKvEntry> attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); |
|
|
|
|
|
|
|
Map<String, Long> subState = new HashMap<>(keys.size()); |
|
|
|
keys.forEach(key -> subState.put(key, 0L)); |
|
|
|
@ -458,6 +458,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
|
|
|
|
TbAttributeSubscriptionScope scope = StringUtils.isEmpty(cmd.getScope()) ? TbAttributeSubscriptionScope.ANY_SCOPE : TbAttributeSubscriptionScope.valueOf(cmd.getScope()); |
|
|
|
|
|
|
|
Lock subLock = new ReentrantLock(); |
|
|
|
TbAttributeSubscription sub = TbAttributeSubscription.builder() |
|
|
|
.serviceId(serviceId) |
|
|
|
.sessionId(sessionId) |
|
|
|
@ -467,9 +468,24 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
.allKeys(false) |
|
|
|
.keyStates(subState) |
|
|
|
.scope(scope) |
|
|
|
.updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) |
|
|
|
.updateConsumer((sessionId, update) -> { |
|
|
|
subLock.lock(); |
|
|
|
try { |
|
|
|
sendWsMsg(sessionId, update); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
}) |
|
|
|
.build(); |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
|
|
|
|
subLock.lock(); |
|
|
|
try{ |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -550,13 +566,13 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
@Override |
|
|
|
public void onSuccess(List<AttributeKvEntry> data) { |
|
|
|
List<TsKvEntry> attributesData = data.stream().map(d -> new BasicTsKvEntry(d.getLastUpdateTs(), d)).collect(Collectors.toList()); |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); |
|
|
|
|
|
|
|
Map<String, Long> subState = new HashMap<>(attributesData.size()); |
|
|
|
attributesData.forEach(v -> subState.put(v.getKey(), v.getTs())); |
|
|
|
|
|
|
|
TbAttributeSubscriptionScope scope = StringUtils.isEmpty(cmd.getScope()) ? TbAttributeSubscriptionScope.ANY_SCOPE : TbAttributeSubscriptionScope.valueOf(cmd.getScope()); |
|
|
|
|
|
|
|
Lock subLock = new ReentrantLock(); |
|
|
|
TbAttributeSubscription sub = TbAttributeSubscription.builder() |
|
|
|
.serviceId(serviceId) |
|
|
|
.sessionId(sessionId) |
|
|
|
@ -565,9 +581,24 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
.entityId(entityId) |
|
|
|
.allKeys(true) |
|
|
|
.keyStates(subState) |
|
|
|
.updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) |
|
|
|
.scope(scope).build(); |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
.updateConsumer((sessionId, update) -> { |
|
|
|
subLock.lock(); |
|
|
|
try { |
|
|
|
sendWsMsg(sessionId, update); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
}) |
|
|
|
.scope(scope) |
|
|
|
.build(); |
|
|
|
|
|
|
|
subLock.lock(); |
|
|
|
try { |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), attributesData)); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -636,20 +667,34 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
FutureCallback<List<TsKvEntry>> callback = new FutureCallback<List<TsKvEntry>>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(List<TsKvEntry> data) { |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); |
|
|
|
Map<String, Long> subState = new HashMap<>(data.size()); |
|
|
|
data.forEach(v -> subState.put(v.getKey(), v.getTs())); |
|
|
|
|
|
|
|
Lock subLock = new ReentrantLock(); |
|
|
|
TbTimeseriesSubscription sub = TbTimeseriesSubscription.builder() |
|
|
|
.serviceId(serviceId) |
|
|
|
.sessionId(sessionId) |
|
|
|
.subscriptionId(cmd.getCmdId()) |
|
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
|
.entityId(entityId) |
|
|
|
.updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) |
|
|
|
.updateConsumer((sessionId, update) -> { |
|
|
|
subLock.lock(); |
|
|
|
try { |
|
|
|
sendWsMsg(sessionId, update); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
}) |
|
|
|
.allKeys(true) |
|
|
|
.keyStates(subState).build(); |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
|
|
|
|
subLock.lock(); |
|
|
|
try { |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -673,21 +718,35 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
return new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(List<TsKvEntry> data) { |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); |
|
|
|
Map<String, Long> subState = new HashMap<>(keys.size()); |
|
|
|
keys.forEach(key -> subState.put(key, startTs)); |
|
|
|
data.forEach(v -> subState.put(v.getKey(), v.getTs())); |
|
|
|
|
|
|
|
Lock subLock = new ReentrantLock(); |
|
|
|
TbTimeseriesSubscription sub = TbTimeseriesSubscription.builder() |
|
|
|
.serviceId(serviceId) |
|
|
|
.sessionId(sessionId) |
|
|
|
.subscriptionId(cmd.getCmdId()) |
|
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
|
.entityId(entityId) |
|
|
|
.updateConsumer(DefaultTelemetryWebSocketService.this::sendWsMsg) |
|
|
|
.updateConsumer((sessionId, update) -> { |
|
|
|
subLock.lock(); |
|
|
|
try { |
|
|
|
sendWsMsg(sessionId, update); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
}) |
|
|
|
.allKeys(false) |
|
|
|
.keyStates(subState).build(); |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
|
|
|
|
subLock.lock(); |
|
|
|
try{ |
|
|
|
oldSubService.addSubscription(sub); |
|
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(cmd.getCmdId(), data)); |
|
|
|
} finally { |
|
|
|
subLock.unlock(); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -793,7 +852,7 @@ public class DefaultTelemetryWebSocketService implements TelemetryWebSocketServi |
|
|
|
|
|
|
|
private void sendWsMsg(TelemetryWebSocketSessionRef sessionRef, int cmdId, Object update) { |
|
|
|
try { |
|
|
|
String msg = jsonMapper.writeValueAsString(update); |
|
|
|
String msg = JacksonUtil.OBJECT_MAPPER.writeValueAsString(update); |
|
|
|
executor.submit(() -> { |
|
|
|
try { |
|
|
|
msgEndpoint.send(sessionRef, cmdId, msg); |
|
|
|
|