|
|
@ -259,13 +259,13 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void sendUpdate(String sessionId, TelemetrySubscriptionUpdate update) { |
|
|
public void sendUpdate(String sessionId, int cmdId, TelemetrySubscriptionUpdate update) { |
|
|
sendUpdate(sessionId, update.getSubscriptionId(), update); |
|
|
doSendUpdate(sessionId, cmdId, update); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
public void sendUpdate(String sessionId, CmdUpdate update) { |
|
|
public void sendUpdate(String sessionId, CmdUpdate update) { |
|
|
sendUpdate(sessionId, update.getCmdId(), update); |
|
|
doSendUpdate(sessionId, update.getCmdId(), update); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
@Override |
|
|
@Override |
|
|
@ -274,7 +274,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
sendUpdate(sessionRef, update); |
|
|
sendUpdate(sessionRef, update); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
private <T> void sendUpdate(String sessionId, int cmdId, T update) { |
|
|
private <T> void doSendUpdate(String sessionId, int cmdId, T update) { |
|
|
WsSessionMetaData md = wsSessionsMap.get(sessionId); |
|
|
WsSessionMetaData md = wsSessionsMap.get(sessionId); |
|
|
if (md != null) { |
|
|
if (md != null) { |
|
|
sendUpdate(md.getSessionRef(), cmdId, update); |
|
|
sendUpdate(md.getSessionRef(), cmdId, update); |
|
|
@ -288,7 +288,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
try { |
|
|
try { |
|
|
msgEndpoint.close(md.getSessionRef(), status); |
|
|
msgEndpoint.close(md.getSessionRef(), status); |
|
|
} catch (IOException e) { |
|
|
} catch (IOException e) { |
|
|
log.warn("[{}] Failed to send session close: {}", sessionId, e); |
|
|
log.warn("[{}] Failed to send session close", sessionId, e); |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
} |
|
|
@ -439,7 +439,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
TbAttributeSubscription sub = TbAttributeSubscription.builder() |
|
|
TbAttributeSubscription sub = TbAttributeSubscription.builder() |
|
|
.serviceId(serviceId) |
|
|
.serviceId(serviceId) |
|
|
.sessionId(sessionId) |
|
|
.sessionId(sessionId) |
|
|
.subscriptionId(cmd.getCmdId()) |
|
|
.subscriptionId(sessionRef.getSessionSubIdSeq().incrementAndGet()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.entityId(entityId) |
|
|
.entityId(entityId) |
|
|
.queryTs(queryTs) |
|
|
.queryTs(queryTs) |
|
|
@ -449,7 +449,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
.updateProcessor((subscription, update) -> { |
|
|
.updateProcessor((subscription, update) -> { |
|
|
subLock.lock(); |
|
|
subLock.lock(); |
|
|
try { |
|
|
try { |
|
|
sendUpdate(subscription.getSessionId(), update); |
|
|
sendUpdate(subscription.getSessionId(), cmd.getCmdId(), update); |
|
|
} finally { |
|
|
} finally { |
|
|
subLock.unlock(); |
|
|
subLock.unlock(); |
|
|
} |
|
|
} |
|
|
@ -545,7 +545,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
TbAttributeSubscription sub = TbAttributeSubscription.builder() |
|
|
TbAttributeSubscription sub = TbAttributeSubscription.builder() |
|
|
.serviceId(serviceId) |
|
|
.serviceId(serviceId) |
|
|
.sessionId(sessionId) |
|
|
.sessionId(sessionId) |
|
|
.subscriptionId(cmd.getCmdId()) |
|
|
.subscriptionId(sessionRef.getSessionSubIdSeq().incrementAndGet()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.entityId(entityId) |
|
|
.entityId(entityId) |
|
|
.queryTs(queryTs) |
|
|
.queryTs(queryTs) |
|
|
@ -554,7 +554,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
.updateProcessor((subscription, update) -> { |
|
|
.updateProcessor((subscription, update) -> { |
|
|
subLock.lock(); |
|
|
subLock.lock(); |
|
|
try { |
|
|
try { |
|
|
sendUpdate(subscription.getSessionId(), update); |
|
|
sendUpdate(subscription.getSessionId(), cmd.getCmdId(), update); |
|
|
} finally { |
|
|
} finally { |
|
|
subLock.unlock(); |
|
|
subLock.unlock(); |
|
|
} |
|
|
} |
|
|
@ -643,13 +643,13 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
TbTimeSeriesSubscription sub = TbTimeSeriesSubscription.builder() |
|
|
TbTimeSeriesSubscription sub = TbTimeSeriesSubscription.builder() |
|
|
.serviceId(serviceId) |
|
|
.serviceId(serviceId) |
|
|
.sessionId(sessionId) |
|
|
.sessionId(sessionId) |
|
|
.subscriptionId(cmd.getCmdId()) |
|
|
.subscriptionId(sessionRef.getSessionSubIdSeq().incrementAndGet()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.entityId(entityId) |
|
|
.entityId(entityId) |
|
|
.updateProcessor((subscription, update) -> { |
|
|
.updateProcessor((subscription, update) -> { |
|
|
subLock.lock(); |
|
|
subLock.lock(); |
|
|
try { |
|
|
try { |
|
|
sendUpdate(subscription.getSessionId(), update); |
|
|
sendUpdate(subscription.getSessionId(), cmd.getCmdId(), update); |
|
|
} finally { |
|
|
} finally { |
|
|
subLock.unlock(); |
|
|
subLock.unlock(); |
|
|
} |
|
|
} |
|
|
@ -698,13 +698,13 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
TbTimeSeriesSubscription sub = TbTimeSeriesSubscription.builder() |
|
|
TbTimeSeriesSubscription sub = TbTimeSeriesSubscription.builder() |
|
|
.serviceId(serviceId) |
|
|
.serviceId(serviceId) |
|
|
.sessionId(sessionId) |
|
|
.sessionId(sessionId) |
|
|
.subscriptionId(cmd.getCmdId()) |
|
|
.subscriptionId(sessionRef.getSessionSubIdSeq().incrementAndGet()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.tenantId(sessionRef.getSecurityCtx().getTenantId()) |
|
|
.entityId(entityId) |
|
|
.entityId(entityId) |
|
|
.updateProcessor((subscription, update) -> { |
|
|
.updateProcessor((subscription, update) -> { |
|
|
subLock.lock(); |
|
|
subLock.lock(); |
|
|
try { |
|
|
try { |
|
|
sendUpdate(subscription.getSessionId(), update); |
|
|
sendUpdate(subscription.getSessionId(), cmd.getCmdId(), update); |
|
|
} finally { |
|
|
} finally { |
|
|
subLock.unlock(); |
|
|
subLock.unlock(); |
|
|
} |
|
|
} |
|
|
@ -836,7 +836,7 @@ public class DefaultWebSocketService implements WebSocketService { |
|
|
try { |
|
|
try { |
|
|
msgEndpoint.sendPing(md.getSessionRef(), currentTime); |
|
|
msgEndpoint.sendPing(md.getSessionRef(), currentTime); |
|
|
} catch (IOException e) { |
|
|
} catch (IOException e) { |
|
|
log.warn("[{}] Failed to send ping: {}", md.getSessionRef().getSessionId(), e); |
|
|
log.warn("[{}] Failed to send ping:", md.getSessionRef().getSessionId(), e); |
|
|
} |
|
|
} |
|
|
})); |
|
|
})); |
|
|
} |
|
|
} |
|
|
|