From e9011eb6469082dd4fc437a62c4b2258e5d30ac0 Mon Sep 17 00:00:00 2001 From: Alexander Skoblikov Date: Wed, 22 Mar 2023 22:54:12 +0100 Subject: [PATCH] CB-3152 allow multiple ws connections in same session (#1534) Co-authored-by: Serge Rider Co-authored-by: dariamarutkina <125263541+dariamarutkina@users.noreply.github.com> --- .../websockets/CBJettyWebSocketManager.java | 35 ++++++++++--------- 1 file changed, 19 insertions(+), 16 deletions(-) diff --git a/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java b/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java index e73700ff82..e6e5498497 100644 --- a/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java +++ b/server/bundles/io.cloudbeaver.server/src/io/cloudbeaver/server/websockets/CBJettyWebSocketManager.java @@ -28,12 +28,14 @@ import org.jkiss.dbeaver.Log; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; public class CBJettyWebSocketManager implements JettyWebSocketCreator { private static final Log log = Log.getLog(CBJettyWebSocketManager.class); - private final Map socketBySessionId = new ConcurrentHashMap<>(); + private final Map> socketBySessionId = new ConcurrentHashMap<>(); private final WebSessionManager webSessionManager; public CBJettyWebSocketManager(WebSessionManager webSessionManager) { @@ -54,12 +56,8 @@ public class CBJettyWebSocketManager implements JettyWebSocketCreator { return null; } var webSessionId = webSession.getSessionId(); - var oldWebSocket = socketBySessionId.get(webSessionId); - if (oldWebSocket != null) { - oldWebSocket.close(); - } var newWebSocket = new CBEventsWebSocket(webSession); - socketBySessionId.put(webSessionId, newWebSocket); + socketBySessionId.computeIfAbsent(webSessionId, key -> new ArrayList<>()).add(newWebSocket); log.info("Websocket created for session: " + webSessionId); return newWebSocket; } @@ -83,8 +81,11 @@ public class CBJettyWebSocketManager implements JettyWebSocketCreator { public void sendPing() { //remove expired sessions socketBySessionId.entrySet() - .removeIf(entry -> - entry.getValue().isNotConnected() || webSessionManager.getSession(entry.getKey()) == null + .removeIf(entry -> { + entry.getValue().removeIf(ws -> !ws.isConnected()); + return webSessionManager.getSession(entry.getKey()) == null || + entry.getValue().isEmpty(); + } ); socketBySessionId.entrySet() @@ -93,14 +94,16 @@ public class CBJettyWebSocketManager implements JettyWebSocketCreator { .forEach( entry -> { var sessionId = entry.getKey(); - var webSocket = entry.getValue(); - try { - webSocket.getRemote().sendPing( - ByteBuffer.wrap("cb-ping".getBytes(StandardCharsets.UTF_8)), - webSocket.getCallback() - ); - } catch (Exception e) { - log.error("Failed to send ping in web socket: " + sessionId); + var webSockets = entry.getValue(); + for (CBEventsWebSocket webSocket : webSockets) { + try { + webSocket.getRemote().sendPing( + ByteBuffer.wrap("cb-ping".getBytes(StandardCharsets.UTF_8)), + webSocket.getCallback() + ); + } catch (Exception e) { + log.error("Failed to send ping in web socket: " + sessionId); + } } } );