@ -13,7 +13,7 @@
* See the License for the specific language governing permissions and
* limitations under the License .
* /
package org.thingsboard.server.service.telemetry ;
package org.thingsboard.server.service.ws ;
import com.fasterxml.jackson.core.JsonProcessingException ;
import com.fasterxml.jackson.databind.ObjectMapper ;
@ -62,14 +62,9 @@ import org.thingsboard.server.service.subscription.TbAttributeSubscriptionScope;
import org.thingsboard.server.service.subscription.TbEntityDataSubscriptionService ;
import org.thingsboard.server.service.subscription.TbLocalSubscriptionService ;
import org.thingsboard.server.service.subscription.TbTimeseriesSubscription ;
import org.thingsboard.server.service.ws.SessionEvent ;
import org.thingsboard.server.service.ws.WebSocketMsgEndpoint ;
import org.thingsboard.server.service.ws.WebSocketSessionRef ;
import org.thingsboard.server.service.ws.WsCmd ;
import org.thingsboard.server.service.ws.WsSessionMetaData ;
import org.thingsboard.server.service.ws.notification.NotificationCommandsHandler ;
import org.thingsboard.server.service.ws.notification.cmd.NotificationCmdsWrapper ;
import org.thingsboard.server.service.ws.telemetry.WebSocketService ;
import org.thingsboard.server.service.ws.notification.cmd.WsCmd ;
import org.thingsboard.server.service.ws.telemetry.cmd.TelemetryPluginCmdsWrapper ;
import org.thingsboard.server.service.ws.telemetry.cmd.v1.AttributesSubscriptionCmd ;
import org.thingsboard.server.service.ws.telemetry.cmd.v1.GetHistoryCmd ;
@ -163,22 +158,22 @@ public class DefaultWebSocketService implements WebSocketService {
pingExecutor . scheduleWithFixedDelay ( this : : sendPing , pingTimeout / NUMBER_OF_PING_ATTEMPTS , pingTimeout / NUMBER_OF_PING_ATTEMPTS , TimeUnit . MILLISECONDS ) ;
telemetryCmdsHandlers = List . of (
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getAttrSubCmds , this : : handleWsAttributesSubscriptionCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getTsSubCmds , this : : handleWsTimeseriesSubscriptionCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getHistoryCmds , this : : handleWsHistoryCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getEntityDataCmds , this : : handleWsEntityDataCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getAlarmDataCmds , this : : handleWsAlarmDataCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getEntityCountCmds , this : : handleWsEntityCountCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getEntityDataUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getAlarmDataUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getAlarmDataUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd ) ,
WsCmdListHandler . of ( TelemetryPluginCmdsWrapper : : getEntityCountUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd )
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getAttrSubCmds , this : : handleWsAttributesSubscriptionCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getTsSubCmds , this : : handleWsTimeseriesSubscriptionCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getHistoryCmds , this : : handleWsHistoryCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getEntityDataCmds , this : : handleWsEntityDataCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getAlarmDataCmds , this : : handleWsAlarmDataCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getEntityCountCmds , this : : handleWsEntityCountCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getEntityDataUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getAlarmDataUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getAlarmDataUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd ) ,
newCmdsHandler ( TelemetryPluginCmdsWrapper : : getEntityCountUnsubscribeCmds , this : : handleWsDataUnsubscribeCmd )
) ;
notificationCmdsHandlers = List . of (
WsCmdHandler . of ( NotificationCmdsWrapper : : getUnreadSubCmd , notificationCmdsHandler : : handleUnreadNotificationsSubCmd ) ,
WsCmdHandler . of ( NotificationCmdsWrapper : : getUnreadCountSubCmd , notificationCmdsHandler : : handleUnreadNotificationsCountSubCmd ) ,
WsCmdHandler . of ( NotificationCmdsWrapper : : getMarkAsReadCmd , notificationCmdsHandler : : handleMarkAsReadCmd ) ,
WsCmdHandler . of ( NotificationCmdsWrapper : : getUnsubCmd , notificationCmdsHandler : : handleUnsubCmd )
newCmdHandler ( NotificationCmdsWrapper : : getUnreadSubCmd , notificationCmdsHandler : : handleUnreadNotificationsSubCmd ) ,
newCmdHandler ( NotificationCmdsWrapper : : getUnreadCountSubCmd , notificationCmdsHandler : : handleUnreadNotificationsCountSubCmd ) ,
newCmdHandler ( NotificationCmdsWrapper : : getMarkAsReadCmd , notificationCmdsHandler : : handleMarkAsReadCmd ) ,
newCmdHandler ( NotificationCmdsWrapper : : getUnsubCmd , notificationCmdsHandler : : handleUnsubCmd )
) ;
}
@ -955,7 +950,17 @@ public class DefaultWebSocketService implements WebSocketService {
return limit = = 0 ? DEFAULT_LIMIT : limit ;
}
@RequiredArgsConstructor ( staticName = "of" )
public static < W , C > WsCmdHandler < W , C > newCmdHandler ( java . util . function . Function < W , C > cmdExtractor ,
BiConsumer < WebSocketSessionRef , C > handler ) {
return new WsCmdHandler < > ( cmdExtractor , handler ) ;
}
public static < W , C > WsCmdListHandler < W , C > newCmdsHandler ( java . util . function . Function < W , List < C > > cmdsExtractor ,
BiConsumer < WebSocketSessionRef , C > handler ) {
return new WsCmdListHandler < > ( cmdsExtractor , handler ) ;
}
@RequiredArgsConstructor
public static class WsCmdHandler < W , C > {
private final java . util . function . Function < W , C > cmdExtractor ;
private final BiConsumer < WebSocketSessionRef , C > handler ;
@ -970,13 +975,13 @@ public class DefaultWebSocketService implements WebSocketService {
}
}
@RequiredArgsConstructor ( staticName = "of" )
@RequiredArgsConstructor
public static class WsCmdListHandler < W , C > {
private final java . util . function . Function < W , List < C > > cmdExtractor ;
private final java . util . function . Function < W , List < C > > cmds Extractor ;
private final BiConsumer < WebSocketSessionRef , C > handler ;
public List < C > extractCmds ( W cmdsWrapper ) {
return cmdExtractor . apply ( cmdsWrapper ) ;
return cmds Extractor . apply ( cmdsWrapper ) ;
}
@SuppressWarnings ( "unchecked" )