|
|
@ -21,8 +21,10 @@ import com.google.common.util.concurrent.FutureCallback; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.Futures; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
import com.google.common.util.concurrent.ListenableFuture; |
|
|
import com.google.common.util.concurrent.MoreExecutors; |
|
|
import com.google.common.util.concurrent.MoreExecutors; |
|
|
|
|
|
import lombok.Getter; |
|
|
import lombok.RequiredArgsConstructor; |
|
|
import lombok.RequiredArgsConstructor; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
import lombok.extern.slf4j.Slf4j; |
|
|
|
|
|
import org.apache.commons.collections.CollectionUtils; |
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
import org.springframework.beans.factory.annotation.Value; |
|
|
import org.springframework.stereotype.Service; |
|
|
import org.springframework.stereotype.Service; |
|
|
import org.springframework.web.socket.CloseStatus; |
|
|
import org.springframework.web.socket.CloseStatus; |
|
|
@ -65,8 +67,7 @@ import org.thingsboard.server.service.subscription.TbLocalSubscriptionService; |
|
|
import org.thingsboard.server.service.subscription.TbTimeSeriesSubscription; |
|
|
import org.thingsboard.server.service.subscription.TbTimeSeriesSubscription; |
|
|
import org.thingsboard.server.service.ws.notification.NotificationCommandsHandler; |
|
|
import org.thingsboard.server.service.ws.notification.NotificationCommandsHandler; |
|
|
import org.thingsboard.server.service.ws.notification.cmd.NotificationCmdsWrapper; |
|
|
import org.thingsboard.server.service.ws.notification.cmd.NotificationCmdsWrapper; |
|
|
import org.thingsboard.server.service.ws.notification.cmd.WsCmd; |
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.TelemetryCmdsWrapper; |
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.WsCmdsWrapper; |
|
|
|
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.v1.AttributesSubscriptionCmd; |
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.v1.AttributesSubscriptionCmd; |
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.v1.GetHistoryCmd; |
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.v1.GetHistoryCmd; |
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.v1.SubscriptionCmd; |
|
|
import org.thingsboard.server.service.ws.telemetry.cmd.v1.SubscriptionCmd; |
|
|
@ -149,7 +150,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
private ScheduledExecutorService pingExecutor; |
|
|
private ScheduledExecutorService pingExecutor; |
|
|
private String serviceId; |
|
|
private String serviceId; |
|
|
|
|
|
|
|
|
private List<WsCmdsHandler<? extends WsCmd>> cmdsHandlers; |
|
|
private List<WsCmdHandler<? extends WsCmd>> cmdsHandlers; |
|
|
|
|
|
|
|
|
@PostConstruct |
|
|
@PostConstruct |
|
|
public void init() { |
|
|
public void init() { |
|
|
@ -160,22 +161,22 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
pingExecutor.scheduleWithFixedDelay(this::sendPing, pingTimeout / NUMBER_OF_PING_ATTEMPTS, pingTimeout / NUMBER_OF_PING_ATTEMPTS, TimeUnit.MILLISECONDS); |
|
|
pingExecutor.scheduleWithFixedDelay(this::sendPing, pingTimeout / NUMBER_OF_PING_ATTEMPTS, pingTimeout / NUMBER_OF_PING_ATTEMPTS, TimeUnit.MILLISECONDS); |
|
|
|
|
|
|
|
|
cmdsHandlers = List.of( |
|
|
cmdsHandlers = List.of( |
|
|
newCmdsHandler(WsCmdsWrapper::getAttrSubCmds, this::handleWsAttributesSubscriptionCmd), |
|
|
newCmdHandler(WsCmdType.ATTRIBUTES, this::handleWsAttributesSubscriptionCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getTsSubCmds, this::handleWsTimeseriesSubscriptionCmd), |
|
|
newCmdHandler(WsCmdType.TIMESERIES, this::handleWsTimeseriesSubscriptionCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getHistoryCmds, this::handleWsHistoryCmd), |
|
|
newCmdHandler(WsCmdType.TIMESERIES_HISTORY, this::handleWsHistoryCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getEntityDataCmds, this::handleWsEntityDataCmd), |
|
|
newCmdHandler(WsCmdType.ENTITY_DATA, this::handleWsEntityDataCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getAlarmDataCmds, this::handleWsAlarmDataCmd), |
|
|
newCmdHandler(WsCmdType.ALARM_DATA, this::handleWsAlarmDataCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getEntityCountCmds, this::handleWsEntityCountCmd), |
|
|
newCmdHandler(WsCmdType.ENTITY_COUNT, this::handleWsEntityCountCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getAlarmCountCmds, this::handleWsAlarmCountCmd), |
|
|
newCmdHandler(WsCmdType.ALARM_COUNT, this::handleWsAlarmCountCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getEntityDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdHandler(WsCmdType.ENTITY_DATA_UNSUBSCRIBE, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getAlarmDataUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdHandler(WsCmdType.ALARM_DATA_UNSUBSCRIBE, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getEntityCountUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdHandler(WsCmdType.ENTITY_COUNT_UNSUBSCRIBE, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdsHandler(WsCmdsWrapper::getAlarmCountUnsubscribeCmds, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdHandler(WsCmdType.ALARM_COUNT_UNSUBSCRIBE, this::handleWsDataUnsubscribeCmd), |
|
|
newCmdHandler(WsCmdsWrapper::getUnreadNotificationsSubCmd, notificationCmdsHandler::handleUnreadNotificationsSubCmd), |
|
|
newCmdHandler(WsCmdType.NOTIFICATIONS, notificationCmdsHandler::handleUnreadNotificationsSubCmd), |
|
|
newCmdHandler(WsCmdsWrapper::getUnreadNotificationsCountSubCmd, notificationCmdsHandler::handleUnreadNotificationsCountSubCmd), |
|
|
newCmdHandler(WsCmdType.NOTIFICATIONS_COUNT, notificationCmdsHandler::handleUnreadNotificationsCountSubCmd), |
|
|
newCmdHandler(WsCmdsWrapper::getMarkNotificationAsReadCmd, notificationCmdsHandler::handleMarkAsReadCmd), |
|
|
newCmdHandler(WsCmdType.MARK_NOTIFICATIONS_AS_READ, notificationCmdsHandler::handleMarkAsReadCmd), |
|
|
newCmdHandler(WsCmdsWrapper::getMarkAllNotificationsAsReadCmd, notificationCmdsHandler::handleMarkAllAsReadCmd), |
|
|
newCmdHandler(WsCmdType.MARK_ALL_NOTIFICATIONS_AS_READ, notificationCmdsHandler::handleMarkAllAsReadCmd), |
|
|
newCmdHandler(WsCmdsWrapper::getNotificationsUnsubCmd, notificationCmdsHandler::handleUnsubCmd) |
|
|
newCmdHandler(WsCmdType.NOTIFICATIONS_UNSUBSCRIBE, notificationCmdsHandler::handleUnsubCmd) |
|
|
); |
|
|
); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@ -217,91 +218,71 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
try { |
|
|
try { |
|
|
|
|
|
WsCommandsWrapper cmdsWrapper; |
|
|
switch (sessionRef.getSessionType()) { |
|
|
switch (sessionRef.getSessionType()) { |
|
|
case GENERAL: |
|
|
case GENERAL: |
|
|
processCmds(sessionRef, msg); |
|
|
cmdsWrapper = JacksonUtil.fromString(msg, WsCommandsWrapper.class); |
|
|
|
|
|
break; |
|
|
|
|
|
case TELEMETRY: |
|
|
|
|
|
cmdsWrapper = JacksonUtil.fromString(msg, TelemetryCmdsWrapper.class).toCommonCmdsWrapper(); |
|
|
break; |
|
|
break; |
|
|
case NOTIFICATIONS: |
|
|
case NOTIFICATIONS: |
|
|
processNotificationCmds(sessionRef, msg); |
|
|
cmdsWrapper = JacksonUtil.fromString(msg, NotificationCmdsWrapper.class).toCommonCmdsWrapper(); |
|
|
break; |
|
|
break; |
|
|
|
|
|
default: |
|
|
|
|
|
throw new IllegalArgumentException("Unknown session type"); |
|
|
} |
|
|
} |
|
|
} catch (IOException e) { |
|
|
processCmds(sessionRef, cmdsWrapper); |
|
|
|
|
|
} catch (Exception e) { |
|
|
log.warn("Failed to decode subscription cmd: {}", e.getMessage(), e); |
|
|
log.warn("Failed to decode subscription cmd: {}", e.getMessage(), e); |
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(UNKNOWN_SUBSCRIPTION_ID, SubscriptionErrorCode.BAD_REQUEST, FAILED_TO_PARSE_WS_COMMAND)); |
|
|
sendWsMsg(sessionRef, new TelemetrySubscriptionUpdate(UNKNOWN_SUBSCRIPTION_ID, SubscriptionErrorCode.BAD_REQUEST, FAILED_TO_PARSE_WS_COMMAND)); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void processCmds(WebSocketSessionRef sessionRef, String msg) throws JsonProcessingException { |
|
|
private void processCmds(WebSocketSessionRef sessionRef, WsCommandsWrapper cmdsWrapper) { |
|
|
WsCmdsWrapper cmdsWrapper = JacksonUtil.fromString(msg, WsCmdsWrapper.class); |
|
|
if (cmdsWrapper == null || CollectionUtils.isEmpty(cmdsWrapper.getCmds())) { |
|
|
processCmds(sessionRef, cmdsWrapper); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void processNotificationCmds(WebSocketSessionRef sessionRef, String msg) throws IOException { |
|
|
|
|
|
NotificationCmdsWrapper cmdsWrapper = JacksonUtil.fromString(msg, NotificationCmdsWrapper.class); |
|
|
|
|
|
processCmds(sessionRef, cmdsWrapper.toCommonCmdsWrapper()); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void processCmds(WebSocketSessionRef sessionRef, WsCmdsWrapper cmdsWrapper) { |
|
|
|
|
|
if (cmdsWrapper == null) { |
|
|
|
|
|
return; |
|
|
return; |
|
|
} |
|
|
} |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
for (WsCmdsHandler<? extends WsCmd> cmdHandler : cmdsHandlers) { |
|
|
if (!validateSessionMetadata(sessionRef, cmdsWrapper.getCmds().get(0).getCmdId(), sessionId)) { |
|
|
List<? extends WsCmd> cmds = cmdHandler.extract(cmdsWrapper); |
|
|
return; |
|
|
if (cmds == null) { |
|
|
} |
|
|
continue; |
|
|
|
|
|
} |
|
|
for (WsCmd cmd : cmdsWrapper.getCmds()) { |
|
|
for (WsCmd cmd : cmds) { |
|
|
log.debug("[{}][{}][{}] Processing cmd: {}", sessionId, cmd.getType(), cmd.getCmdId(), cmd); |
|
|
if (validateSessionMetadata(sessionRef, cmd.getCmdId(), sessionId)) { |
|
|
try { |
|
|
try { |
|
|
getCmdHandler(cmd.getType()).handle(sessionRef, cmd); |
|
|
cmdHandler.handle(sessionRef, cmd); |
|
|
} catch (Exception e) { |
|
|
} catch (Exception e) { |
|
|
log.error("[sessionId: {}, tenantId: {}, userId: {}] Failed to handle WS cmd: {}", sessionId, |
|
|
log.error("[sessionId: {}, tenantId: {}, userId: {}] Failed to handle WS cmd: {}", sessionId, |
|
|
sessionRef.getSecurityCtx().getTenantId(), sessionRef.getSecurityCtx().getId(), cmd, e); |
|
|
sessionRef.getSecurityCtx().getTenantId(), sessionRef.getSecurityCtx().getId(), cmd, e); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void handleWsEntityDataCmd(WebSocketSessionRef sessionRef, EntityDataCmd cmd) { |
|
|
private void handleWsEntityDataCmd(WebSocketSessionRef sessionRef, EntityDataCmd cmd) { |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
|
|
|
log.debug("[{}] Processing: {}", sessionId, cmd); |
|
|
|
|
|
|
|
|
|
|
|
if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void handleWsEntityCountCmd(WebSocketSessionRef sessionRef, EntityCountCmd cmd) { |
|
|
private void handleWsEntityCountCmd(WebSocketSessionRef sessionRef, EntityCountCmd cmd) { |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
|
|
|
log.debug("[{}] Processing: {}", sessionId, cmd); |
|
|
|
|
|
|
|
|
|
|
|
if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void handleWsAlarmDataCmd(WebSocketSessionRef sessionRef, AlarmDataCmd cmd) { |
|
|
private void handleWsAlarmDataCmd(WebSocketSessionRef sessionRef, AlarmDataCmd cmd) { |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
|
|
|
log.debug("[{}] Processing: {}", sessionId, cmd); |
|
|
|
|
|
|
|
|
|
|
|
if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void handleWsDataUnsubscribeCmd(WebSocketSessionRef sessionRef, UnsubscribeCmd cmd) { |
|
|
private void handleWsDataUnsubscribeCmd(WebSocketSessionRef sessionRef, UnsubscribeCmd cmd) { |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
|
|
|
log.debug("[{}] Processing: {}", sessionId, cmd); |
|
|
|
|
|
entityDataSubService.cancelSubscription(sessionRef.getSessionId(), cmd); |
|
|
entityDataSubService.cancelSubscription(sessionRef.getSessionId(), cmd); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void handleWsAlarmCountCmd(WebSocketSessionRef sessionRef, AlarmCountCmd cmd) { |
|
|
private void handleWsAlarmCountCmd(WebSocketSessionRef sessionRef, AlarmCountCmd cmd) { |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
if (validateCmd(sessionRef, cmd)) { |
|
|
log.debug("[{}] Processing: {}", sessionId, cmd); |
|
|
|
|
|
|
|
|
|
|
|
if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
|
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
entityDataSubService.handleCmd(sessionRef, cmd); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
@ -447,8 +428,6 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
log.debug("[{}] Processing: {}", sessionId, cmd); |
|
|
|
|
|
|
|
|
|
|
|
if (cmd.isUnsubscribe()) { |
|
|
if (cmd.isUnsubscribe()) { |
|
|
unsubscribe(sessionRef, cmd, sessionId); |
|
|
unsubscribe(sessionRef, cmd, sessionId); |
|
|
} else if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
} else if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
@ -533,18 +512,15 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private void handleWsHistoryCmd(WebSocketSessionRef sessionRef, GetHistoryCmd cmd) { |
|
|
private void handleWsHistoryCmd(WebSocketSessionRef sessionRef, GetHistoryCmd cmd) { |
|
|
if (cmd.getEntityId() == null || cmd.getEntityId().isEmpty() || cmd.getEntityType() == null || cmd.getEntityType().isEmpty()) { |
|
|
if (!validateCmd(sessionRef, cmd, () -> { |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
if (cmd.getEntityId() == null || cmd.getEntityId().isEmpty() || cmd.getEntityType() == null || cmd.getEntityType().isEmpty()) { |
|
|
"Device id is empty!"); |
|
|
throw new IllegalArgumentException("Device id is empty!"); |
|
|
sendWsMsg(sessionRef, update); |
|
|
} |
|
|
return; |
|
|
if (cmd.getKeys() == null || cmd.getKeys().isEmpty()) { |
|
|
} |
|
|
throw new IllegalArgumentException("Keys are empty!"); |
|
|
if (cmd.getKeys() == null || cmd.getKeys().isEmpty()) { |
|
|
} |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
})) return; |
|
|
"Keys are empty!"); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
return; |
|
|
|
|
|
} |
|
|
|
|
|
EntityId entityId = EntityIdFactory.getByTypeAndId(cmd.getEntityType(), cmd.getEntityId()); |
|
|
EntityId entityId = EntityIdFactory.getByTypeAndId(cmd.getEntityType(), cmd.getEntityId()); |
|
|
List<String> keys = new ArrayList<>(getKeys(cmd).orElse(Collections.emptySet())); |
|
|
List<String> keys = new ArrayList<>(getKeys(cmd).orElse(Collections.emptySet())); |
|
|
List<ReadTsKvQuery> queries = keys.stream().map(key -> new BaseReadTsKvQuery(key, cmd.getStartTs(), cmd.getEndTs(), cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()))) |
|
|
List<ReadTsKvQuery> queries = keys.stream().map(key -> new BaseReadTsKvQuery(key, cmd.getStartTs(), cmd.getEndTs(), cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()))) |
|
|
@ -621,9 +597,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
@Override |
|
|
@Override |
|
|
public void onFailure(Throwable e) { |
|
|
public void onFailure(Throwable e) { |
|
|
log.error(FAILED_TO_FETCH_ATTRIBUTES, e); |
|
|
log.error(FAILED_TO_FETCH_ATTRIBUTES, e); |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR, |
|
|
sendError(sessionRef, cmd.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR, FAILED_TO_FETCH_ATTRIBUTES); |
|
|
FAILED_TO_FETCH_ATTRIBUTES); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
} |
|
|
} |
|
|
}; |
|
|
}; |
|
|
|
|
|
|
|
|
@ -641,8 +615,6 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
String sessionId = sessionRef.getSessionId(); |
|
|
log.debug("[{}] Processing: {}", sessionId, cmd); |
|
|
|
|
|
|
|
|
|
|
|
if (cmd.isUnsubscribe()) { |
|
|
if (cmd.isUnsubscribe()) { |
|
|
unsubscribe(sessionRef, cmd, sessionId); |
|
|
unsubscribe(sessionRef, cmd, sessionId); |
|
|
} else if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
} else if (validateSubscriptionCmd(sessionRef, cmd)) { |
|
|
@ -781,9 +753,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
} else { |
|
|
} else { |
|
|
log.info(FAILED_TO_FETCH_DATA, e); |
|
|
log.info(FAILED_TO_FETCH_DATA, e); |
|
|
} |
|
|
} |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR, |
|
|
sendError(sessionRef, cmd.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR, FAILED_TO_FETCH_DATA); |
|
|
FAILED_TO_FETCH_DATA); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
} |
|
|
} |
|
|
}; |
|
|
}; |
|
|
} |
|
|
} |
|
|
@ -797,86 +767,73 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, EntityDataCmd cmd) { |
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, EntityDataCmd cmd) { |
|
|
if (cmd.getCmdId() < 0) { |
|
|
return validateCmd(sessionRef, cmd, () -> { |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
if (cmd.getQuery() == null && !cmd.hasAnyCmd()) { |
|
|
"Cmd id is negative value!"); |
|
|
throw new IllegalArgumentException("Query is empty!"); |
|
|
sendWsMsg(sessionRef, update); |
|
|
} |
|
|
return false; |
|
|
}); |
|
|
} else if (cmd.getQuery() == null && !cmd.hasAnyCmd()) { |
|
|
|
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
|
|
|
"Query is empty!"); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
return false; |
|
|
|
|
|
} |
|
|
|
|
|
return true; |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, EntityCountCmd cmd) { |
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, EntityCountCmd cmd) { |
|
|
if (cmd.getCmdId() < 0) { |
|
|
return validateCmd(sessionRef, cmd, () -> { |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
if (cmd.getQuery() == null) { |
|
|
"Cmd id is negative value!"); |
|
|
throw new IllegalArgumentException("Query is empty!"); |
|
|
sendWsMsg(sessionRef, update); |
|
|
} |
|
|
return false; |
|
|
}); |
|
|
} else if (cmd.getQuery() == null) { |
|
|
|
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, "Query is empty!"); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
return false; |
|
|
|
|
|
} |
|
|
|
|
|
return true; |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, AlarmDataCmd cmd) { |
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, AlarmDataCmd cmd) { |
|
|
if (cmd.getCmdId() < 0) { |
|
|
return validateCmd(sessionRef, cmd, () -> { |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
if (cmd.getQuery() == null) { |
|
|
"Cmd id is negative value!"); |
|
|
throw new IllegalArgumentException("Query is empty!"); |
|
|
sendWsMsg(sessionRef, update); |
|
|
} |
|
|
return false; |
|
|
}); |
|
|
} else if (cmd.getQuery() == null) { |
|
|
|
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
|
|
|
"Query is empty!"); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
return false; |
|
|
|
|
|
} |
|
|
|
|
|
return true; |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, SubscriptionCmd cmd) { |
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, SubscriptionCmd cmd) { |
|
|
if (cmd.getEntityId() == null || cmd.getEntityId().isEmpty()) { |
|
|
return validateCmd(sessionRef, cmd, () -> { |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
if (cmd.getEntityId() == null || cmd.getEntityId().isEmpty()) { |
|
|
"Device id is empty!"); |
|
|
throw new IllegalArgumentException("Device id is empty!"); |
|
|
sendWsMsg(sessionRef, update); |
|
|
} |
|
|
return false; |
|
|
}); |
|
|
} |
|
|
|
|
|
return true; |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private boolean validateSessionMetadata(WebSocketSessionRef sessionRef, SubscriptionCmd cmd, String sessionId) { |
|
|
|
|
|
return validateSessionMetadata(sessionRef, cmd.getCmdId(), sessionId); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean validateSessionMetadata(WebSocketSessionRef sessionRef, int cmdId, String sessionId) { |
|
|
private boolean validateSessionMetadata(WebSocketSessionRef sessionRef, int cmdId, String sessionId) { |
|
|
WsSessionMetaData sessionMD = wsSessionsMap.get(sessionId); |
|
|
WsSessionMetaData sessionMD = wsSessionsMap.get(sessionId); |
|
|
if (sessionMD == null) { |
|
|
if (sessionMD == null) { |
|
|
log.warn("[{}] Session meta data not found. ", sessionId); |
|
|
log.warn("[{}] Session meta data not found. ", sessionId); |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmdId, SubscriptionErrorCode.INTERNAL_ERROR, |
|
|
sendError(sessionRef, cmdId, SubscriptionErrorCode.INTERNAL_ERROR, SESSION_META_DATA_NOT_FOUND); |
|
|
SESSION_META_DATA_NOT_FOUND); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
return false; |
|
|
return false; |
|
|
} else { |
|
|
} else { |
|
|
return true; |
|
|
return true; |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private boolean validateSubscriptionCmd(WebSocketSessionRef sessionRef, AlarmCountCmd cmd) { |
|
|
private boolean validateCmd(WebSocketSessionRef sessionRef, WsCmd cmd) { |
|
|
|
|
|
return validateCmd(sessionRef, cmd, null); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private <C extends WsCmd> boolean validateCmd(WebSocketSessionRef sessionRef, C cmd, Runnable validator) { |
|
|
if (cmd.getCmdId() < 0) { |
|
|
if (cmd.getCmdId() < 0) { |
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, |
|
|
sendError(sessionRef, cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, "Cmd id is negative value!"); |
|
|
"Cmd id is negative value!"); |
|
|
return false; |
|
|
sendWsMsg(sessionRef, update); |
|
|
} |
|
|
|
|
|
try { |
|
|
|
|
|
if (validator != null) { |
|
|
|
|
|
validator.run(); |
|
|
|
|
|
} |
|
|
|
|
|
} catch (Exception e) { |
|
|
|
|
|
sendError(sessionRef, cmd.getCmdId(), SubscriptionErrorCode.BAD_REQUEST, e.getMessage()); |
|
|
return false; |
|
|
return false; |
|
|
} |
|
|
} |
|
|
return true; |
|
|
return true; |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void sendError(WebSocketSessionRef sessionRef, int subId, SubscriptionErrorCode errorCode, String errorMsg) { |
|
|
|
|
|
TelemetrySubscriptionUpdate update = new TelemetrySubscriptionUpdate(subId, errorCode, errorMsg); |
|
|
|
|
|
sendWsMsg(sessionRef, update); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
private void sendWsMsg(WebSocketSessionRef sessionRef, EntityDataUpdate update) { |
|
|
private void sendWsMsg(WebSocketSessionRef sessionRef, EntityDataUpdate update) { |
|
|
sendWsMsg(sessionRef, update.getCmdId(), update); |
|
|
sendWsMsg(sessionRef, update.getCmdId(), update); |
|
|
} |
|
|
} |
|
|
@ -1034,39 +991,29 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
.map(TenantProfile::getDefaultProfileConfiguration).orElse(null); |
|
|
.map(TenantProfile::getDefaultProfileConfiguration).orElse(null); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
public WsCmdHandler<? extends WsCmd> getCmdHandler(WsCmdType cmdType) { |
|
|
public static <C extends WsCmd> WsCmdHandler<C> newCmdHandler(java.util.function.Function<WsCmdsWrapper, C> cmdExtractor, |
|
|
for (WsCmdHandler<? extends WsCmd> cmdHandler : cmdsHandlers) { |
|
|
BiConsumer<WebSocketSessionRef, C> handler) { |
|
|
if (cmdHandler.getCmdType() == cmdType) { |
|
|
return new WsCmdHandler<>(cmdExtractor, handler); |
|
|
return cmdHandler; |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
throw new IllegalArgumentException("Unknown command type " + cmdType); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
public static <C extends WsCmd> WsCmdsHandler<C> newCmdsHandler(java.util.function.Function<WsCmdsWrapper, List<C>> cmdsExtractor, |
|
|
public static <C extends WsCmd> WsCmdHandler<C> newCmdHandler(WsCmdType cmdType, BiConsumer<WebSocketSessionRef, C> handler) { |
|
|
BiConsumer<WebSocketSessionRef, C> handler) { |
|
|
return new WsCmdHandler<>(cmdType, handler); |
|
|
return new WsCmdsHandler<>(cmdsExtractor, handler); |
|
|
|
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@RequiredArgsConstructor |
|
|
@RequiredArgsConstructor |
|
|
public static class WsCmdsHandler<C extends WsCmd> { |
|
|
@Getter |
|
|
private final java.util.function.Function<WsCmdsWrapper, List<C>> cmdsExtractor; |
|
|
@SuppressWarnings("unchecked") |
|
|
|
|
|
public static class WsCmdHandler<C extends WsCmd> { |
|
|
|
|
|
private final WsCmdType cmdType; |
|
|
protected final BiConsumer<WebSocketSessionRef, C> handler; |
|
|
protected final BiConsumer<WebSocketSessionRef, C> handler; |
|
|
|
|
|
|
|
|
public List<C> extract(WsCmdsWrapper cmdsWrapper) { |
|
|
public void handle(WebSocketSessionRef sessionRef, WsCmd cmd) { |
|
|
return cmdsExtractor.apply(cmdsWrapper); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings("unchecked") |
|
|
|
|
|
public void handle(WebSocketSessionRef sessionRef, Object cmd) { |
|
|
|
|
|
handler.accept(sessionRef, (C) cmd); |
|
|
handler.accept(sessionRef, (C) cmd); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
public static class WsCmdHandler<C extends WsCmd> extends WsCmdsHandler<C> { |
|
|
|
|
|
public WsCmdHandler(java.util.function.Function<WsCmdsWrapper, C> cmdExtractor, BiConsumer<WebSocketSessionRef, C> handler) { |
|
|
|
|
|
super(cmdsWrapper -> { |
|
|
|
|
|
C cmd = cmdExtractor.apply(cmdsWrapper); |
|
|
|
|
|
return cmd != null ? List.of(cmd) : null; |
|
|
|
|
|
}, handler); |
|
|
|
|
|
} |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
} |
|
|
} |
|
|
|