diff --git a/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocket.java b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocket.java new file mode 100644 index 00000000..707fc342 --- /dev/null +++ b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocket.java @@ -0,0 +1,229 @@ +/* + * The contents of this file are subject to the terms of the Common Development and + * Distribution License (the License). You may not use this file except in compliance with the + * License. + * + * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the + * specific language governing permission and limitations under the License. + * + * When distributing Covered Software, include this CDDL Header Notice in each file and include + * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL + * Header, with the fields enclosed by brackets [] replaced by your own identifying + * information: "Portions copyright [year] [name of copyright owner]". + * + * Copyright 2015-2016 ForgeRock AS. + * Portions Copyrighted 2026 3A Systems, LLC + */ +package org.forgerock.openicf.framework.server.jetty; + +import java.io.IOException; +import java.nio.ByteBuffer; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.eclipse.jetty.websocket.api.BatchMode; +import org.eclipse.jetty.websocket.api.Frame; +import org.eclipse.jetty.websocket.api.Session; +import org.eclipse.jetty.websocket.api.StatusCode; +import org.eclipse.jetty.websocket.api.WebSocketFrameListener; +import org.eclipse.jetty.websocket.api.WebSocketListener; +import org.eclipse.jetty.websocket.api.WebSocketPingPongListener; +import org.forgerock.openicf.common.protobuf.RPCMessages; +import org.forgerock.openicf.framework.remote.ConnectionPrincipal; +import org.forgerock.openicf.framework.remote.rpc.RemoteOperationContext; +import org.forgerock.openicf.framework.remote.rpc.WebSocketConnectionHolder; +import org.forgerock.util.Utils; +import org.forgerock.util.promise.Promises; +import org.identityconnectors.common.logging.Log; +import org.identityconnectors.framework.common.exceptions.ConnectorIOException; + +/** + * The Jetty endpoint of a single WebSocket connection. + *

+ * One instance is created per connection by {@link OpenICFWebSocketCreator}, + * while the {@link ConnectionPrincipal} it belongs to is cached per principal + * name and shared by every connection of that name. Keeping the session, the + * {@link WebSocketConnectionHolder} and the close state here (instead of on + * the shared principal) lets concurrent and successive connections of one + * principal coexist: each connection delivers its own close event and cleans + * up its own resources. + */ +public class OpenICFWebSocket implements + WebSocketPingPongListener, WebSocketListener, WebSocketFrameListener { + + private static final Log logger = Log.getLog(OpenICFWebSocket.class); + + private final ConnectionPrincipal principal; + + // Written on the Jetty connect thread, read on the send-executor thread + // and by isOperational() callers. + private volatile Session session; + + private final AtomicBoolean closed = new AtomicBoolean(false); + + // Written on the handshake-processing pool thread, read by other message + // threads via getRemoteConnectionContext()/isHandHooked(). + private volatile RemoteOperationContext context = null; + + // Single send thread per connection: frames must leave in submission + // order (the peer drops e.g. an operation response that overtakes the + // handshake response). Jetty's RemoteEndpoint is thread-safe, but + // concurrent blocking sends may reach the wire in any order. The executor + // is shut down when this connection closes. + private final ExecutorService sendExecutor = Executors.newSingleThreadExecutor( + Utils.newThreadFactory(null, "OpenICF Jetty WebSocket Send %d", true)); + + public OpenICFWebSocket(final ConnectionPrincipal principal) { + this.principal = principal; + } + + Session getSession() { + return session; + } + + @Override + public void onWebSocketConnect(Session session) { + WebSocketPingPongListener.super.onWebSocketConnect(session); + this.session = session; + principal.getOperationMessageListener().onConnect(adapter); + } + + @Override + public void onWebSocketClose(int statusCode, String reason) { + if (!closed.compareAndSet(false, true)) { + return; + } + try { + principal.getOperationMessageListener().onClose(adapter, statusCode, reason); + } finally { + // Notifies the holder's close listeners so the + // WebSocketConnectionGroup drops this connection. + adapter.close(); + sendExecutor.shutdown(); + } + } + + @Override + public void onWebSocketError(Throwable t) { + logger.warn(t, "onError"); + principal.getOperationMessageListener().onError(t); + } + + @Override + public void onWebSocketPing(ByteBuffer buffer) { + byte[] b = new byte[buffer.remaining()]; + buffer.get(b); + principal.getOperationMessageListener().onPing(adapter, b); + } + + @Override + public void onWebSocketPong(ByteBuffer buffer) { + byte[] b = new byte[buffer.remaining()]; + buffer.get(b); + principal.getOperationMessageListener().onPong(adapter, b); + } + + @Override + public void onWebSocketBinary(byte[] payload, int offset, int len) { + logger.ok("onBinaryMessage({0})", null != payload ? payload.length : 0); + principal.getOperationMessageListener().onMessage(adapter, payload); + } + + @Override + public void onWebSocketText(String message) { + logger.ok("onTextMessage({0})", message); + principal.getOperationMessageListener().onMessage(adapter, message); + } + + @Override + public void onWebSocketFrame(Frame frame) { + logger.ok("onWebSocketFrame({0})", frame); + } + + private final WebSocketConnectionHolder adapter = new WebSocketConnectionHolder() { + + @Override + protected void handshake(RPCMessages.HandshakeMessage message) { + context = principal.handshake(this, message); + } + + @Override + public boolean isOperational() { + return null != getSession() && getSession().isOpen(); + } + + @Override + public RemoteOperationContext getRemoteConnectionContext() { + return context; + } + + @Override + public Future sendBytes(byte[] data) { + return submitSend(() -> getSession().getRemote().sendBytes(ByteBuffer.wrap(data))); + } + + @Override + public Future sendString(String data) { + return submitSend(() -> getSession().getRemote().sendString(data)); + } + + private Future submitSend(FrameWrite frame) { + if (isOperational()) { + try { + return sendExecutor.submit(() -> { + try { + frame.write(); + } catch (IOException e) { + throw new RuntimeException(e); + } + }); + } catch (RejectedExecutionException e) { + // The connection was closed and the executor shut down. + } + } + return Promises.newExceptionPromise(new ConnectorIOException( + "Socket is not connected.")); + } + + @Override + public void sendPing(byte[] applicationData) throws Exception { + if (isOperational()) { + getSession().getRemote().sendPing(ByteBuffer.wrap(applicationData)); + if (getSession().getRemote().getBatchMode() == BatchMode.ON) { + getSession().getRemote().flush(); + } + } else { + throw new ConnectorIOException("Socket is not connected."); + } + } + + @Override + public void sendPong(byte[] applicationData) throws Exception { + if (isOperational()) { + getSession().getRemote().sendPong(ByteBuffer.wrap(applicationData)); + if (getSession().getRemote().getBatchMode() == BatchMode.ON) { + getSession().getRemote().flush(); + } + } else { + throw new ConnectorIOException("Socket is not connected."); + } + } + + @Override + protected void tryClose() { + final Session current = getSession(); + if (null != current && current.isOpen()) { + current.close(StatusCode.NORMAL, "Shutdown"); + } + } + + }; + + @FunctionalInterface + private interface FrameWrite { + void write() throws IOException; + } +} diff --git a/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketCreator.java b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketCreator.java index c8c16229..b6f6e18e 100644 --- a/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketCreator.java +++ b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketCreator.java @@ -12,15 +12,17 @@ * information: "Portions copyright [year] [name of copyright owner]". * * Copyright 2015-2016 ForgeRock AS. - * Portions copyright 2025 3A Systems LLC. + * Portions copyright 2025-2026 3A Systems LLC. */ package org.forgerock.openicf.framework.server.jetty; +import java.io.Closeable; import java.io.IOException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import javax.security.auth.callback.NameCallback; @@ -39,7 +41,7 @@ import org.forgerock.openicf.framework.remote.rpc.OperationMessageListener; import org.forgerock.openicf.framework.remote.rpc.WebSocketConnectionGroup; -public class OpenICFWebSocketCreator implements JettyWebSocketCreator { +public class OpenICFWebSocketCreator implements JettyWebSocketCreator, Closeable { private static final Logger logger = Log.getLogger(OpenICFWebSocketCreator.class); @@ -52,6 +54,12 @@ public class OpenICFWebSocketCreator implements JettyWebSocketCreator { private Authenticator authenticator; + private final ScheduledFuture groupHealthChecker; + + // Set by close(); createWebSocket()/authenticate() must not mint new + // principals over a framework that is being released. + private volatile boolean closed = false; + public OpenICFWebSocketCreator(final ConnectorFramework connectorFramework, final ScheduledExecutorService executorService) { @@ -84,8 +92,8 @@ public void authenticate(JettyServerUpgradeRequest request, JettyServerUpgradeRe this.authenticator = authenticator; } - //Will be cancelled when executorService is shut down - executorService.scheduleWithFixedDelay(new Runnable() { + //Cancelled in close(); also dies when executorService is shut down + groupHealthChecker = executorService.scheduleWithFixedDelay(new Runnable() { public void run() { for (WebSocketConnectionGroup e : globalConnectionGroups.values()) { if (!e.checkIsActive()) { @@ -99,19 +107,61 @@ public void run() { @Override public Object createWebSocket(JettyServerUpgradeRequest request, JettyServerUpgradeResponse response) { + if (closed) { + if (!response.isCommitted()) { + try { + response.sendError(HttpServletResponse.SC_SERVICE_UNAVAILABLE, + "OpenICF connector server is shut down"); + } catch (IOException e) { + // + } + } + return null; + } + if (request.getSubProtocols().contains(RemoteWSFrameworkConnectionInfo.OPENICF_PROTOCOL)) { response.setAcceptedSubProtocol(RemoteWSFrameworkConnectionInfo.OPENICF_PROTOCOL); } ConnectionPrincipal connectionPrincipal = authenticate(request, response); if (null != connectionPrincipal) { - return connectionPrincipal; + // The principal is cached and shared by every connection of this + // name; the endpoint holds the per-connection state. + return new OpenICFWebSocket(connectionPrincipal); } else if (!response.isCommitted()) { unauthorized(response, "Unknown Principal"); } return null; } + /** + * Ends the lifecycle of every cached principal: shuts their connection + * groups down and notifies the principals' close listeners. Invoked from + * {@link OpenICFWebSocketServletBase#destroy()}. + */ + @Override + public void close() { + closed = true; + groupHealthChecker.cancel(false); + // One pass over the groups: every cached principal leaves every group + // before the group is dropped, so pending requests of all principals + // are cancelled (removing the group inside a per-principal loop would + // hide it from the remaining principals). + for (WebSocketConnectionGroup group : globalConnectionGroups.values()) { + for (ConnectionPrincipal principal : principalCache.values()) { + // We should gracefully shut down the group + group.principalIsShuttingDown(principal); + } + if (!group.isOperational()) { + globalConnectionGroups.remove(group.getRemoteSessionId()); + } + } + for (ConnectionPrincipal principal : principalCache.values()) { + principal.close(); + } + principalCache.clear(); + } + protected void unauthorized(JettyServerUpgradeResponse response, String message) { try { response.sendError( @@ -123,17 +173,15 @@ protected void unauthorized(JettyServerUpgradeResponse response, String message) } public ConnectionPrincipal authenticate(JettyServerUpgradeRequest request, JettyServerUpgradeResponse response) { + if (closed) { + return null; + } NameCallback callback = new NameCallback("OpenICF user:>"); authenticator.authenticate(request, response, callback); if (StringUtil.isNotBlank(callback.getName())) { - ConnectionPrincipal connectionPrincipal = principalCache.get(callback.getName()); - if (connectionPrincipal == null) { - principalCache.putIfAbsent(callback.getName(), new SinglePrincipal(callback.getName(), listener, - connectorFramework, globalConnectionGroups)); - return principalCache.get(callback.getName()); - } else { - return connectionPrincipal; - } + return principalCache.computeIfAbsent(callback.getName(), + name -> new SinglePrincipal(name, listener, connectorFramework, + globalConnectionGroups)); } return null; } diff --git a/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketServletBase.java b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketServletBase.java index 2388b249..62fe0c90 100644 --- a/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketServletBase.java +++ b/OpenICF-java-framework/connector-server-jetty/src/main/java/org/forgerock/openicf/framework/server/jetty/OpenICFWebSocketServletBase.java @@ -12,7 +12,7 @@ * information: "Portions copyright [year] [name of copyright owner]". * * Copyright 2015-2016 ForgeRock AS. - * Portions copyright 2025 3A Systems LLC. + * Portions copyright 2025-2026 3A Systems LLC. */ package org.forgerock.openicf.framework.server.jetty; @@ -56,6 +56,7 @@ public class OpenICFWebSocketServletBase extends JettyWebSocketServlet { private ReferenceCountedObject.Reference connectorFramework = null; private ScheduledExecutorService executorService = null; + private OpenICFWebSocketCreator websocketCreator = null; public OpenICFWebSocketServletBase() { } @@ -73,6 +74,17 @@ public OpenICFWebSocketServletBase(final ReferenceCountedObject implements - WebSocketPingPongListener, WebSocketListener, WebSocketFrameListener { +/** + * The authenticated principal of one or more WebSocket connections. + *

+ * {@link OpenICFWebSocketCreator} caches one instance per principal name and + * closes it from {@link OpenICFWebSocketCreator#close()}. The per-connection + * state lives in {@link OpenICFWebSocket}, which is created for every + * accepted connection. + */ +public class SinglePrincipal extends ConnectionPrincipal { final String name; final ConnectorFramework connectorFramework; @@ -75,180 +63,12 @@ public RemoteOperationContext handshake( } protected void doClose() { - sendExecutor.shutdown(); + // Per-connection resources are released by OpenICFWebSocket when the + // connection closes; the principal itself holds none. } - @Override protected void onNewWebSocketConnectionGroup(final WebSocketConnectionGroup connectionGroup) { connectorFramework.getServerManager(getName()).addWebSocketConnectionGroup(connectionGroup); } - - @Override - public void onWebSocketPing(ByteBuffer buffer) { - byte[] b = new byte[buffer.remaining()]; - buffer.get(b); - getConnectionPrincipal().getOperationMessageListener().onPing(adapter, b); - } - - @Override - public void onWebSocketPong(ByteBuffer buffer) { - byte[] b = new byte[buffer.remaining()]; - buffer.get(b); - getConnectionPrincipal().getOperationMessageListener().onPong(adapter, b); - } - - @Override - public void onWebSocketClose(int statusCode, String reason) { - if (hasCloseBeenCalled) { - return; - } - hasCloseBeenCalled = true; - getConnectionPrincipal().getOperationMessageListener().onClose(adapter, - statusCode, reason); - } - - Session session; - - @Override - public void onWebSocketConnect(Session session) { - WebSocketPingPongListener.super.onWebSocketConnect(session); - this.session = session; - getConnectionPrincipal().getOperationMessageListener().onConnect(adapter); - } - - Session getSession() { - return this.session; - } - - @Override - public void onWebSocketError(Throwable t) { - logger.debug("onError:", t); - getConnectionPrincipal().getOperationMessageListener().onError(t); - } - - @Override - public void onWebSocketBinary(byte[] payload, int offset, int len) { - logger.debug("onBinaryMessage('" + (null != payload ? payload.length : 0) + "')"); - getConnectionPrincipal().getOperationMessageListener().onMessage(adapter, payload); - } - - @Override - public void onWebSocketText(String message) { - logger.debug("onTextMessage('" + message + "')"); - getConnectionPrincipal().getOperationMessageListener().onMessage(adapter, message); - } - - @Override - public void onWebSocketFrame(Frame frame) { - logger.debug("onWebSocketFrame('" + frame + "')"); - } - - private static final Logger logger = Log.getLogger(SinglePrincipal.class); - private boolean hasCloseBeenCalled = false; - - // Written on the handshake-processing pool thread, read by other message - // threads via getRemoteConnectionContext()/isHandHooked(). - private volatile RemoteOperationContext context = null; - - // Single send thread per principal: frames must leave in submission order - // (the peer drops e.g. an operation response that overtakes the handshake - // response). Jetty's RemoteEndpoint is thread-safe, but concurrent - // blocking sends may reach the wire in any order. This instance is cached - // by OpenICFWebSocketCreator and serves every connection of this - // principal name, so the executor must survive onWebSocketClose; it is - // shut down in doClose(). The thread is a daemon because cached - // principals are not closed on servlet destroy. - private final ExecutorService sendExecutor = Executors.newSingleThreadExecutor( - Utils.newThreadFactory(null, "OpenICF Jetty WebSocket Send %d", true)); - - private final WebSocketConnectionHolder adapter = new WebSocketConnectionHolder() { - - protected void handshake(RPCMessages.HandshakeMessage message) { - context = getConnectionPrincipal().handshake(this, message); - } - - public boolean isOperational() { - return getSession().isOpen(); - } - - public RemoteOperationContext getRemoteConnectionContext() { - return context; - } - - public Future sendBytes(byte[] data) { - if (isOperational()) { - try { - return sendExecutor.submit(() -> { - try { - getSession().getRemote().sendBytes(ByteBuffer.wrap(data)); - } catch (IOException e) { - throw new RuntimeException(e); - } - }); - } catch (RejectedExecutionException e) { - // The principal was closed and doClose() shut the - // executor down. - return Promises.newExceptionPromise(new ConnectorIOException( - "Socket is not connected.")); - } - } else { - return Promises.newExceptionPromise(new ConnectorIOException( - "Socket is not connected.")); - } - } - - public Future sendString(String data) { - if (isOperational()) { - try { - return sendExecutor.submit(() -> { - try { - getSession().getRemote().sendString(data); - } catch (IOException e) { - throw new RuntimeException(e); - } - }); - } catch (RejectedExecutionException e) { - return Promises.newExceptionPromise(new ConnectorIOException( - "Socket is not connected.")); - } - } else { - return Promises.newExceptionPromise(new ConnectorIOException( - "Socket is not connected.")); - } - } - - public void sendPing(byte[] applicationData) throws Exception { - if (isOperational()) { - getSession().getRemote().sendPing(ByteBuffer.wrap(applicationData)); - if (getSession().getRemote().getBatchMode() == BatchMode.ON) { - getSession().getRemote().flush(); - } - } else { - throw new ConnectorIOException("Socket is not connected."); - } - } - - public void sendPong(byte[] applicationData) throws Exception { - if (isOperational()) { - getSession().getRemote().sendPong(ByteBuffer.wrap(applicationData)); - if (getSession().getRemote().getBatchMode() == BatchMode.ON) { - getSession().getRemote().flush(); - } - } else { - throw new ConnectorIOException("Socket is not connected."); - } - } - - protected void tryClose() { - getSession().close(StatusCode.NORMAL, "TEST003"); - } - - }; - - protected ConnectionPrincipal getConnectionPrincipal() { - return this; - } - - } diff --git a/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/PrincipalLifecycleTest.java b/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/PrincipalLifecycleTest.java new file mode 100644 index 00000000..0e16174c --- /dev/null +++ b/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/PrincipalLifecycleTest.java @@ -0,0 +1,117 @@ +/* + * The contents of this file are subject to the terms of the Common Development and + * Distribution License (the License). You may not use this file except in compliance with the + * License. + * + * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the + * specific language governing permission and limitations under the License. + * + * When distributing Covered Software, include this CDDL Header Notice in each file and include + * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL + * Header, with the fields enclosed by brackets [] replaced by your own identifying + * information: "Portions copyright [year] [name of copyright owner]". + * + * Copyright 2026 3A Systems, LLC. + */ +package org.forgerock.openicf.framework.server.jetty; + +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.Method; +import java.lang.reflect.Proxy; +import java.util.Collections; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.atomic.AtomicInteger; + +import org.eclipse.jetty.websocket.server.JettyServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.JettyServerUpgradeResponse; +import org.forgerock.openicf.framework.CloseListener; +import org.forgerock.openicf.framework.remote.ConnectionPrincipal; +import org.forgerock.openicf.framework.remote.rpc.OperationMessageListener; +import org.testng.Assert; +import org.testng.annotations.Test; + +/** + * OpenICFWebSocketCreator must cache one principal per name, hand out a fresh + * endpoint per connection, and end the lifecycle of the cached principals + * from close() (OpenIdentityPlatform/OpenICF#112). + */ +public class PrincipalLifecycleTest { + + private static OperationMessageListener noopListener() { + return (OperationMessageListener) Proxy.newProxyInstance( + PrincipalLifecycleTest.class.getClassLoader(), + new Class[] { OperationMessageListener.class }, + new InvocationHandler() { + public Object invoke(Object p, Method m, Object[] a) { + return null; + } + }); + } + + private static JettyServerUpgradeRequest upgradeRequest() { + return (JettyServerUpgradeRequest) Proxy.newProxyInstance( + PrincipalLifecycleTest.class.getClassLoader(), + new Class[] { JettyServerUpgradeRequest.class }, + new InvocationHandler() { + public Object invoke(Object p, Method m, Object[] a) { + if ("getSubProtocols".equals(m.getName())) { + return Collections.emptyList(); + } + return null; + } + }); + } + + private static JettyServerUpgradeResponse upgradeResponse() { + return (JettyServerUpgradeResponse) Proxy.newProxyInstance( + PrincipalLifecycleTest.class.getClassLoader(), + new Class[] { JettyServerUpgradeResponse.class }, + new InvocationHandler() { + public Object invoke(Object p, Method m, Object[] a) { + if ("isCommitted".equals(m.getName())) { + return Boolean.FALSE; + } + return null; + } + }); + } + + @Test(timeOut = 30000) + @SuppressWarnings("unchecked") + public void testCreatorCachesPrincipalAndClosesItOnClose() throws Exception { + ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(1); + try { + OpenICFWebSocketCreator creator = + new OpenICFWebSocketCreator(null, noopListener(), null, scheduler); + + Object first = creator.createWebSocket(upgradeRequest(), upgradeResponse()); + Object second = creator.createWebSocket(upgradeRequest(), upgradeResponse()); + + Assert.assertTrue(first instanceof OpenICFWebSocket, + "creator must hand out a per-connection endpoint"); + Assert.assertNotSame(first, second, + "every connection must get its own endpoint"); + + ConnectionPrincipal principal = creator.authenticate(upgradeRequest(), upgradeResponse()); + Assert.assertSame(creator.authenticate(upgradeRequest(), upgradeResponse()), principal, + "the principal must be cached per name"); + + final AtomicInteger closed = new AtomicInteger(); + ((ConnectionPrincipal) principal) + .addCloseListener(new CloseListener() { + public void onClosed(SinglePrincipal source) { + closed.incrementAndGet(); + } + }); + + creator.close(); + + Assert.assertEquals(closed.get(), 1, + "close() must end the lifecycle of the cached principal"); + Assert.assertTrue(creator.principalCache.isEmpty(), + "close() must drop the cached principals"); + } finally { + scheduler.shutdownNow(); + } + } +} diff --git a/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/ReconnectSendTest.java b/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/ReconnectSendTest.java index 69769286..5c9c6eca 100644 --- a/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/ReconnectSendTest.java +++ b/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/ReconnectSendTest.java @@ -26,6 +26,7 @@ import org.eclipse.jetty.websocket.api.RemoteEndpoint; import org.eclipse.jetty.websocket.api.Session; +import org.forgerock.openicf.common.protobuf.RPCMessages; import org.forgerock.openicf.framework.remote.rpc.OperationMessageListener; import org.forgerock.openicf.framework.remote.rpc.WebSocketConnectionGroup; import org.forgerock.openicf.framework.remote.rpc.WebSocketConnectionHolder; @@ -33,22 +34,23 @@ import org.testng.annotations.Test; /** - * Does a SinglePrincipal that has seen one connection closed still send on a - * subsequent connection? OpenICFWebSocketCreator caches SinglePrincipal per - * principal name, so the same instance serves every connection of that name. + * OpenICFWebSocketCreator caches one SinglePrincipal per principal name and + * creates a fresh OpenICFWebSocket endpoint for every connection. A close on + * one connection must neither break sends on a later connection of the same + * principal nor swallow the close events of later connections. */ public class ReconnectSendTest { - private static final AtomicInteger SENT = new AtomicInteger(); + private final AtomicInteger sent = new AtomicInteger(); - private static Session newSession() { + private Session newSession() { final RemoteEndpoint remote = (RemoteEndpoint) Proxy.newProxyInstance( ReconnectSendTest.class.getClassLoader(), new Class[] { RemoteEndpoint.class }, new InvocationHandler() { public Object invoke(Object p, Method m, Object[] a) { if ("sendBytes".equals(m.getName())) { - SENT.incrementAndGet(); + sent.incrementAndGet(); } return null; } @@ -70,9 +72,10 @@ public Object invoke(Object p, Method m, Object[] a) { } @Test(timeOut = 30000) - public void testSendWorksAfterReconnectOnCachedPrincipal() throws Exception { + public void testSendAndCloseWorkPerConnectionOnCachedPrincipal() throws Exception { final AtomicReference holder = new AtomicReference(); + final AtomicInteger closes = new AtomicInteger(); OperationMessageListener capturing = (OperationMessageListener) Proxy.newProxyInstance( ReconnectSendTest.class.getClassLoader(), @@ -82,6 +85,9 @@ public Object invoke(Object p, Method m, Object[] a) { if ("onConnect".equals(m.getName())) { holder.set((WebSocketConnectionHolder) a[0]); } + if ("onClose".equals(m.getName())) { + closes.incrementAndGet(); + } return null; } }); @@ -90,19 +96,79 @@ public Object invoke(Object p, Method m, Object[] a) { new ConcurrentHashMap()); // --- connection #1 --- - principal.onWebSocketConnect(newSession()); - Future first = holder.get().sendBytes(new byte[] { 1, 2, 3 }); - first.get(5, TimeUnit.SECONDS); - Assert.assertEquals(SENT.get(), 1, "first connection should have sent"); + OpenICFWebSocket first = new OpenICFWebSocket(principal); + first.onWebSocketConnect(newSession()); + Future firstSend = holder.get().sendBytes(new byte[] { 1, 2, 3 }); + firstSend.get(5, TimeUnit.SECONDS); + Assert.assertEquals(sent.get(), 1, "first connection should have sent"); + + first.onWebSocketClose(1000, "client went away"); + Assert.assertEquals(closes.get(), 1, "first close must reach the listener"); + + // duplicate close events of the same connection stay suppressed + first.onWebSocketClose(1000, "duplicate"); + Assert.assertEquals(closes.get(), 1); + + // --- connection #2 on the SAME cached principal --- + OpenICFWebSocket second = new OpenICFWebSocket(principal); + second.onWebSocketConnect(newSession()); + Future secondSend = holder.get().sendBytes(new byte[] { 4, 5, 6 }); + secondSend.get(5, TimeUnit.SECONDS); + Assert.assertEquals(sent.get(), 2, "send after reconnect must reach the wire"); + + second.onWebSocketClose(1000, "client went away again"); + Assert.assertEquals(closes.get(), 2, + "close of a later connection must not be swallowed by an earlier one"); + } + + /** + * The connection group learns about a departed connection only through + * the holder's close listeners (OpenIdentityPlatform/OpenICF#112): before + * the fix onWebSocketClose never called adapter.close(), so every + * reconnect left a stale duplicate holder in the group. + */ + @Test(timeOut = 30000) + public void testGroupDropsHolderOnClose() throws Exception { + final AtomicReference holder = + new AtomicReference(); + + OperationMessageListener capturing = (OperationMessageListener) Proxy.newProxyInstance( + ReconnectSendTest.class.getClassLoader(), + new Class[] { OperationMessageListener.class }, + new InvocationHandler() { + public Object invoke(Object p, Method m, Object[] a) { + if ("onConnect".equals(m.getName())) { + holder.set((WebSocketConnectionHolder) a[0]); + } + return null; + } + }); + + SinglePrincipal principal = new SinglePrincipal("anonymous", capturing, null, + new ConcurrentHashMap()); + WebSocketConnectionGroup group = new WebSocketConnectionGroup("session-1"); + RPCMessages.HandshakeMessage handshake = + RPCMessages.HandshakeMessage.newBuilder().setSessionId("session-1").build(); + + // --- connection #1 joins and leaves the group --- + OpenICFWebSocket first = new OpenICFWebSocket(principal); + first.onWebSocketConnect(newSession()); + group.handshake(principal, holder.get(), handshake); + Assert.assertTrue(group.isOperational(), "handshake must register the holder"); - principal.onWebSocketClose(1000, "client went away"); + first.onWebSocketClose(1000, "client went away"); + Assert.assertFalse(group.isOperational(), + "the group must drop the holder of a closed connection"); - // --- connection #2 on the SAME cached principal instance --- - principal.onWebSocketConnect(newSession()); - Future second = holder.get().sendBytes(new byte[] { 4, 5, 6 }); - second.get(5, TimeUnit.SECONDS); + // --- reconnect: no stale holder of connection #1 may keep the group + // alive after connection #2 leaves as well --- + OpenICFWebSocket second = new OpenICFWebSocket(principal); + second.onWebSocketConnect(newSession()); + group.handshake(principal, holder.get(), handshake); + Assert.assertTrue(group.isOperational(), "reconnect must register the new holder"); - Assert.assertEquals(SENT.get(), 2, - "send after reconnect must reach the wire"); + second.onWebSocketClose(1000, "client went away again"); + Assert.assertFalse(group.isOperational(), + "no duplicate holder may survive a reconnect cycle"); } } diff --git a/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/ServletLifecycleTest.java b/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/ServletLifecycleTest.java new file mode 100644 index 00000000..63bd6982 --- /dev/null +++ b/OpenICF-java-framework/connector-server-jetty/src/test/java/org/forgerock/openicf/framework/server/jetty/ServletLifecycleTest.java @@ -0,0 +1,142 @@ +/* + * The contents of this file are subject to the terms of the Common Development and + * Distribution License (the License). You may not use this file except in compliance with the + * License. + * + * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the + * specific language governing permission and limitations under the License. + * + * When distributing Covered Software, include this CDDL Header Notice in each file and include + * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL + * Header, with the fields enclosed by brackets [] replaced by your own identifying + * information: "Portions copyright [year] [name of copyright owner]". + * + * Copyright 2026 3A Systems, LLC. + */ +package org.forgerock.openicf.framework.server.jetty; + +import java.lang.reflect.InvocationHandler; +import java.lang.reflect.Method; +import java.lang.reflect.Proxy; +import java.util.Collections; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +import org.eclipse.jetty.websocket.server.JettyServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.JettyServerUpgradeResponse; +import org.eclipse.jetty.websocket.server.JettyWebSocketServletFactory; +import org.forgerock.openicf.framework.CloseListener; +import org.forgerock.openicf.framework.ConnectorFramework; +import org.forgerock.openicf.framework.ConnectorFrameworkFactory; +import org.forgerock.openicf.framework.remote.ConnectionPrincipal; +import org.testng.Assert; +import org.testng.annotations.Test; + +/** + * OpenICFWebSocketServletBase.destroy() must end the lifecycle of everything + * the servlet created lazily for itself: close the websocket creator (and + * with it the cached principals), shut the private scheduler down and release + * the lazily acquired ConnectorFramework + * (OpenIdentityPlatform/OpenICF#112). + */ +public class ServletLifecycleTest { + + static class TestServlet extends OpenICFWebSocketServletBase { + + private final ConnectorFrameworkFactory factory = new ConnectorFrameworkFactory(); + + volatile ConnectorFramework framework; + + @Override + protected ConnectorFrameworkFactory getConnectorFrameworkFactory() { + return factory; + } + + @Override + protected void configure(ConnectorFramework connectorFramework) { + framework = connectorFramework; + } + } + + private static JettyServerUpgradeRequest upgradeRequest() { + return (JettyServerUpgradeRequest) Proxy.newProxyInstance( + ServletLifecycleTest.class.getClassLoader(), + new Class[] { JettyServerUpgradeRequest.class }, + new InvocationHandler() { + public Object invoke(Object p, Method m, Object[] a) { + if ("getSubProtocols".equals(m.getName())) { + return Collections.emptyList(); + } + return null; + } + }); + } + + private static JettyServerUpgradeResponse upgradeResponse() { + return (JettyServerUpgradeResponse) Proxy.newProxyInstance( + ServletLifecycleTest.class.getClassLoader(), + new Class[] { JettyServerUpgradeResponse.class }, + new InvocationHandler() { + public Object invoke(Object p, Method m, Object[] a) { + if ("isCommitted".equals(m.getName())) { + return Boolean.FALSE; + } + return null; + } + }); + } + + @Test(timeOut = 30000) + @SuppressWarnings("unchecked") + public void testDestroyReleasesLazilyAcquiredResources() throws Exception { + TestServlet servlet = new TestServlet(); + + final AtomicReference creatorRef = + new AtomicReference(); + JettyWebSocketServletFactory factory = + (JettyWebSocketServletFactory) Proxy.newProxyInstance( + ServletLifecycleTest.class.getClassLoader(), + new Class[] { JettyWebSocketServletFactory.class }, + new InvocationHandler() { + public Object invoke(Object p, Method m, Object[] a) { + if ("setCreator".equals(m.getName())) { + creatorRef.set((OpenICFWebSocketCreator) a[0]); + } + return null; + } + }); + + // The no-arg servlet acquires its framework and scheduler lazily. + servlet.configure(factory); + + OpenICFWebSocketCreator creator = creatorRef.get(); + Assert.assertNotNull(creator, "configure() must install the creator"); + Assert.assertNotNull(servlet.framework, "the lazy path must configure the framework"); + Assert.assertTrue(servlet.framework.isRunning()); + + ScheduledExecutorService scheduler = servlet.getExecutorService(); + Assert.assertFalse(scheduler.isShutdown()); + + ConnectionPrincipal principal = creator.authenticate(upgradeRequest(), upgradeResponse()); + Assert.assertNotNull(principal); + final AtomicInteger closed = new AtomicInteger(); + ((ConnectionPrincipal) principal) + .addCloseListener(new CloseListener() { + public void onClosed(SinglePrincipal source) { + closed.incrementAndGet(); + } + }); + + servlet.destroy(); + + Assert.assertEquals(closed.get(), 1, + "destroy() must close the creator and with it the cached principals"); + Assert.assertNull(creator.createWebSocket(upgradeRequest(), upgradeResponse()), + "a closed creator must not accept new connections"); + Assert.assertTrue(scheduler.isShutdown(), + "destroy() must shut the private scheduler down"); + Assert.assertFalse(servlet.framework.isRunning(), + "destroy() must release the lazily acquired framework"); + } +}