This is an automated email from the ASF dual-hosted git repository. bonampak pushed a commit to branch feature/jakarta-jetty-upgrade in repository https://gitbox.apache.org/repos/asf/knox.git
commit 6e86455858f730b0e7bcf17bd3a6879d9216ca62 Author: bonampak <[email protected]> AuthorDate: Fri Apr 24 16:18:26 2026 +0200 KNOX-3238: fix websocket tests --- .../websockets/GatewayWebsocketHandler.java | 2 +- .../webshell/WebshellWebsocketAdapterTest.java | 17 ++-- .../websockets/AbstractWebSocketHandler.java | 100 ++++++++++++++++++ .../knox/gateway/websockets/BadBackendTest.java | 6 +- .../apache/knox/gateway/websockets/BadSocket.java | 38 +++---- .../apache/knox/gateway/websockets/BadUrlTest.java | 4 +- .../gateway/websockets/BigEchoSocketHandler.java | 36 ++++--- .../gateway/websockets/ConnectionDroppedTest.java | 6 +- .../apache/knox/gateway/websockets/EchoSocket.java | 64 ++++++------ .../websockets/GatewayWebsocketHandlerTest.java | 113 ++++++--------------- .../gateway/websockets/MessageFailureTest.java | 8 +- .../gateway/websockets/WebsocketEchoHandler.java | 32 +++--- .../gateway/websockets/WebsocketEchoTestBase.java | 8 +- .../WebsocketMultipleConnectionTest.java | 4 +- .../WebsocketServerInitiatedMessageTest.java | 51 ++++------ .../WebsocketServerInitiatedPingTest.java | 53 ++++------ 16 files changed, 283 insertions(+), 259 deletions(-) diff --git a/gateway-server/src/main/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandler.java b/gateway-server/src/main/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandler.java index 3e3d0ce85..2c9a3630b 100644 --- a/gateway-server/src/main/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandler.java +++ b/gateway-server/src/main/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandler.java @@ -159,7 +159,7 @@ public class GatewayWebsocketHandler extends Handler.Wrapper { // 3. Reconstruct the java.net.URI for Knox's internal routing methods final URI requestURI = HttpURI.build(httpURI).scheme(wsScheme).toURI(); - // Now Knox's regex will work perfectly! + // Now Knox's regex will work if (isWebshellRequest(requestURI)) { return handleWebshellRequest(req); // Note: Update handleWebshellRequest to accept ServerUpgradeRequest } diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/webshell/WebshellWebsocketAdapterTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/webshell/WebshellWebsocketAdapterTest.java index 61e4446ff..84fbe0434 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/webshell/WebshellWebsocketAdapterTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/webshell/WebshellWebsocketAdapterTest.java @@ -27,9 +27,10 @@ import org.apache.knox.gateway.i18n.messages.MessagesFactory; import org.apache.knox.gateway.provider.federation.jwt.JWTMessages; import org.apache.knox.gateway.websockets.JWTValidator; import org.apache.knox.gateway.websockets.WebsocketLogMessages; +import org.easymock.Capture; import org.easymock.EasyMock; import org.easymock.EasyMockSupport; -import org.eclipse.jetty.websocket.api.RemoteEndpoint; +import org.eclipse.jetty.websocket.api.Callback; import org.eclipse.jetty.websocket.api.Session; import org.junit.Rule; import org.junit.Test; @@ -99,14 +100,18 @@ public class WebshellWebsocketAdapterTest extends EasyMockSupport { EasyMock.expect(connectionInfo.getInputStream()).andReturn(inputStreamEmpty).times(1); Session session = EasyMock.createNiceMock(Session.class); - RemoteEndpoint remote = EasyMock.createNiceMock(RemoteEndpoint.class); + Capture<Callback> callbackCapture = EasyMock.newCapture(); + session.sendText(EasyMock.eq(bashOuput), EasyMock.capture(callbackCapture)); + EasyMock.expectLastCall().andAnswer(() -> { + callbackCapture.getValue().succeed(); + return null; + }).anyTimes(); // send to client - EasyMock.expect(session.getRemote()).andReturn(remote).anyTimes(); - remote.sendString(bashOuput); + // clean up EasyMock.expect(session.isOpen()).andReturn(true).anyTimes(); session.close(); - EasyMock.replay(session, remote); + EasyMock.replay(session); connectionInfo.disconnect(); PowerMock.replay(connectionInfo, ConnectionInfo.class); @@ -114,7 +119,7 @@ public class WebshellWebsocketAdapterTest extends EasyMockSupport { ExecutorService pool = Executors.newFixedThreadPool(10); AtomicInteger concurrentWebshells = new AtomicInteger(0); WebshellWebSocketAdapter webshellWebSocketAdapter = new WebshellWebSocketAdapter(pool, gatewayConfig, jwtValidator, concurrentWebshells); - webshellWebSocketAdapter.onWebSocketConnect(session); + webshellWebSocketAdapter.onWebSocketOpen(session); verifyAll(); } diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/AbstractWebSocketHandler.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/AbstractWebSocketHandler.java new file mode 100644 index 000000000..7a14dde29 --- /dev/null +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/AbstractWebSocketHandler.java @@ -0,0 +1,100 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.knox.gateway.websockets; + +import org.eclipse.jetty.server.Handler; +import org.eclipse.jetty.server.Request; +import org.eclipse.jetty.server.Response; +import org.eclipse.jetty.server.Server; +import org.eclipse.jetty.util.Callback; +import org.eclipse.jetty.websocket.server.ServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.ServerUpgradeResponse; +import org.eclipse.jetty.websocket.server.ServerWebSocketContainer; +import org.eclipse.jetty.websocket.server.WebSocketCreator; +import org.eclipse.jetty.websocket.server.WebSocketUpgradeHandler; + +/** + * A base class for Jetty 12 Native WebSocket handlers that abstracts away + * the container initialization and upgrade handler lifecycle. + */ +public abstract class AbstractWebSocketHandler extends Handler.Wrapper implements WebSocketCreator { + + private WebSocketUpgradeHandler wsHandler; + + @Override + protected void doStart() throws Exception { + Server server = getServer(); + if (server == null) { + throw new IllegalStateException("WebSocketHandler must be attached to a Server before starting"); + } + + // Ensure the container exists + ServerWebSocketContainer container = ServerWebSocketContainer.ensure(server, null); + + // Let the subclass apply policies (max sizes, timeouts, etc.) + configure(container); + + // Use the exposed mapping spec + container.addMapping(getMappingSpec(), this); + + // Create, wire, and start the internal UpgradeHandler + this.wsHandler = new WebSocketUpgradeHandler(container); + this.wsHandler.setHandler(getHandler()); + this.wsHandler.start(); + + super.doStart(); + } + + @Override + protected void doStop() throws Exception { + if (this.wsHandler != null) { + this.wsHandler.stop(); + } + super.doStop(); + } + + @Override + public boolean handle(Request request, Response response, Callback callback) throws Exception { + return wsHandler.handle(request, response, callback); + } + + /** + * Defines the URL mapping for this WebSocket handler. + * Defaults to catching all requests. Subclasses can override this to specify exact paths. + * * @return The Jetty mapping spec (e.g., "/ws/*", "regex|^/.*", etc.) + */ + protected String getMappingSpec() { + return "regex|^/.*"; + } + + /** + * Configure the ServerWebSocketContainer policies. + * @param container The Jetty 12 WebSocket container + */ + protected abstract void configure(ServerWebSocketContainer container); + + /** + * Return the Jetty 12 Session.Listener (the WebSocket) for the given request. + * @param req The upgrade request + * @param resp The upgrade response + * @param callback The callback to signal success/failure of the upgrade + * @return The WebSocket instance + */ + @Override + public abstract Object createWebSocket(ServerUpgradeRequest req, ServerUpgradeResponse resp, Callback callback); +} diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadBackendTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadBackendTest.java index 5f9058cfa..33784b5b1 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadBackendTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadBackendTest.java @@ -32,6 +32,7 @@ import jakarta.websocket.ContainerProvider; import jakarta.websocket.WebSocketContainer; import java.net.URI; import java.util.Locale; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -81,8 +82,11 @@ public class BadBackendTest { proxy.addConnector(proxyConnector); /* start Knox with WebsocketAdapter to test */ + final ExecutorService pool = Executors.newFixedThreadPool(10); + final URI backendURI = new URI(BAD_BACKEND); + final BigEchoSocketHandler wsHandler = new BigEchoSocketHandler( - new ProxyWebSocketAdapter(new URI(BAD_BACKEND), Executors.newFixedThreadPool(10), gatewayConfig)); + () -> new ProxyWebSocketAdapter(backendURI, pool, gatewayConfig)); ContextHandler context = new ContextHandler(); context.setContextPath("/"); diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadSocket.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadSocket.java index 0744f8c21..d459044e8 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadSocket.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadSocket.java @@ -17,40 +17,34 @@ */ package org.apache.knox.gateway.websockets; -import org.eclipse.jetty.io.RuntimeIOException; -import org.eclipse.jetty.util.BufferUtil; -import org.eclipse.jetty.websocket.api.BatchMode; -import org.eclipse.jetty.websocket.api.RemoteEndpoint; +import java.nio.ByteBuffer; +import org.eclipse.jetty.websocket.api.Callback; import org.eclipse.jetty.websocket.api.Session; -import org.eclipse.jetty.websocket.api.WebSocketAdapter; - -import java.io.IOException; /** * Simulate a bad socket. * * @since 0.10 */ -class BadSocket extends WebSocketAdapter { +class BadSocket extends Session.Listener.AbstractAutoDemanding { + + private Session session; + @Override - public void onWebSocketConnect(final Session session) { + public void onWebSocketOpen(final Session session) { + super.onWebSocketOpen(session); + this.session = session; } @Override - public void onWebSocketBinary(byte[] payload, int offset, int len) { - if (isNotConnected()) { + public void onWebSocketBinary(ByteBuffer payload, Callback callback) { + if (session == null || !session.isOpen()) { + callback.fail(new IllegalStateException("Session closed")); return; } - try { - RemoteEndpoint remote = getRemote(); - remote.sendBytes(BufferUtil.toBuffer(payload, offset, len), null); - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } - } catch (IOException x) { - throw new RuntimeIOException(x); - } + // Echo the binary payload back to the client and complete the callback + session.sendBinary(payload, callback); } @Override @@ -60,10 +54,10 @@ class BadSocket extends WebSocketAdapter { @Override public void onWebSocketText(String message) { - if (isNotConnected()) { + if (session == null || !session.isOpen()) { return; } // Throw an exception on purpose throw new RuntimeException("Simulating bad connection ..."); } -} +} \ No newline at end of file diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadUrlTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadUrlTest.java index d32ec1c01..5582508e8 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadUrlTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BadUrlTest.java @@ -43,7 +43,7 @@ import org.easymock.EasyMock; import org.eclipse.jetty.server.Server; import org.eclipse.jetty.server.ServerConnector; import org.eclipse.jetty.server.handler.ContextHandler; -import org.eclipse.jetty.server.handler.HandlerCollection; +import org.eclipse.jetty.server.handler.ContextHandlerCollection; import org.junit.AfterClass; import org.junit.Assert; import org.junit.BeforeClass; @@ -154,7 +154,7 @@ public class BadUrlTest { gatewayServer.addConnector(connector); /* workaround so we can add our handler later at runtime */ - HandlerCollection handlers = new HandlerCollection(true); + ContextHandlerCollection handlers = new ContextHandlerCollection(true); /* add some initial handlers */ ContextHandler context = new ContextHandler(); diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BigEchoSocketHandler.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BigEchoSocketHandler.java index 58ecec05b..900d63a63 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BigEchoSocketHandler.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/BigEchoSocketHandler.java @@ -17,33 +17,31 @@ */ package org.apache.knox.gateway.websockets; -import org.eclipse.jetty.websocket.api.WebSocketAdapter; -import org.eclipse.jetty.websocket.server.WebSocketHandler; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeRequest; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeResponse; -import org.eclipse.jetty.websocket.servlet.WebSocketCreator; -import org.eclipse.jetty.websocket.servlet.WebSocketServletFactory; +import org.eclipse.jetty.util.Callback; +import org.eclipse.jetty.websocket.api.Session; +import org.eclipse.jetty.websocket.server.ServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.ServerUpgradeResponse; +import org.eclipse.jetty.websocket.server.ServerWebSocketContainer; + +import java.util.function.Supplier; /** - * A Mock websocket handler that just Echos messages + * A websocket handler that just Echos messages */ -class BigEchoSocketHandler extends WebSocketHandler implements WebSocketCreator { - private final WebSocketAdapter socket; +class BigEchoSocketHandler extends AbstractWebSocketHandler { + private final Supplier<Session.Listener> socketSupplier; - BigEchoSocketHandler(final WebSocketAdapter socket) { - this.socket = socket; + BigEchoSocketHandler(final Supplier<Session.Listener> socketSupplier) { + this.socketSupplier = socketSupplier; } @Override - public void configure(WebSocketServletFactory factory) { - factory.getPolicy().setMaxTextMessageSize(66000); - factory.getPolicy().setMaxTextMessageBufferSize(66000); - factory.setCreator(this); + protected void configure(ServerWebSocketContainer container) { + container.setMaxTextMessageSize(66000); } @Override - public Object createWebSocket(ServletUpgradeRequest req, - ServletUpgradeResponse resp) { - return socket; + public Object createWebSocket(ServerUpgradeRequest req, ServerUpgradeResponse resp, Callback callback) { + return socketSupplier.get(); } -} +} \ No newline at end of file diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/ConnectionDroppedTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/ConnectionDroppedTest.java index 25d3a0851..c462aa044 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/ConnectionDroppedTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/ConnectionDroppedTest.java @@ -31,6 +31,7 @@ import jakarta.websocket.WebSocketContainer; import java.io.IOException; import java.net.URI; import java.util.Locale; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -95,7 +96,7 @@ public class ConnectionDroppedTest { /* start backend with Echo socket */ final BigEchoSocketHandler wsHandler = new BigEchoSocketHandler( - new BadSocket()); + BadSocket::new); ContextHandler context = new ContextHandler(); context.setContextPath("/"); @@ -120,8 +121,9 @@ public class ConnectionDroppedTest { proxy.addConnector(proxyConnector); /* start Knox with WebsocketAdapter to test */ + final ExecutorService pool = Executors.newFixedThreadPool(10); final BigEchoSocketHandler wsHandler = new BigEchoSocketHandler( - new ProxyWebSocketAdapter(serverUri, Executors.newFixedThreadPool(10), gatewayConfig)); + () -> new ProxyWebSocketAdapter(serverUri, pool, gatewayConfig)); ContextHandler context = new ContextHandler(); context.setContextPath("/"); diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/EchoSocket.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/EchoSocket.java index 3e025e7e2..ca8212732 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/EchoSocket.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/EchoSocket.java @@ -17,55 +17,51 @@ */ package org.apache.knox.gateway.websockets; -import java.io.IOException; - -import org.eclipse.jetty.io.RuntimeIOException; -import org.eclipse.jetty.util.BufferUtil; -import org.eclipse.jetty.websocket.api.BatchMode; -import org.eclipse.jetty.websocket.api.RemoteEndpoint; -import org.eclipse.jetty.websocket.api.WebSocketAdapter; +import java.nio.ByteBuffer; +import org.eclipse.jetty.websocket.api.Callback; +import org.eclipse.jetty.websocket.api.Session; /** * A simple Echo socket */ -public class EchoSocket extends WebSocketAdapter { +public class EchoSocket extends Session.Listener.AbstractAutoDemanding { - @Override - public void onWebSocketBinary(byte[] payload, int offset, int len) { - if (isNotConnected()) { - return; - } + private Session session; - try { - RemoteEndpoint remote = getRemote(); - remote.sendBytes(BufferUtil.toBuffer(payload, offset, len), null); - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } - } catch (IOException x) { - throw new RuntimeIOException(x); - } + @Override + public void onWebSocketOpen(Session session) { + super.onWebSocketOpen(session); + this.session = session; } + /** + * Jetty 12 Native API uses ByteBuffer for binary frames. + * The callback must be completed (which happens automatically by passing it to sendBinary) + * to signal Jetty to read the next frame. + */ @Override - public void onWebSocketError(Throwable cause) { - throw new RuntimeException(cause); + public void onWebSocketBinary(ByteBuffer payload, Callback callback) { + if (session == null || !session.isOpen()) { + callback.fail(new IllegalStateException("Session closed")); + return; + } + + // Echo the binary payload back to the client + session.sendBinary(payload, callback); } @Override public void onWebSocketText(String message) { - if (isNotConnected()) { + if (session == null || !session.isOpen()) { return; } - try { - RemoteEndpoint remote = getRemote(); - remote.sendString(message, null); - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } - } catch (IOException x) { - throw new RuntimeIOException(x); - } + // Echo the text message back to the client + session.sendText(message, Callback.NOOP); + } + + @Override + public void onWebSocketError(Throwable cause) { + throw new RuntimeException(cause); } } diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandlerTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandlerTest.java index 47331b807..df4aed6b3 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandlerTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/GatewayWebsocketHandlerTest.java @@ -29,10 +29,10 @@ import org.apache.knox.gateway.provider.federation.jwt.JWTMessages; import org.apache.knox.gateway.services.GatewayServices; import org.apache.knox.gateway.webshell.WebshellWebSocketAdapter; import org.easymock.EasyMock; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeRequest; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeResponse; -import javax.servlet.http.HttpServletRequest; -import javax.servlet.http.HttpServletResponse; +import org.eclipse.jetty.http.HttpFields; +import org.eclipse.jetty.http.HttpURI; +import org.eclipse.jetty.websocket.server.ServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.ServerUpgradeResponse; import org.junit.Assert; import org.junit.BeforeClass; import org.junit.Rule; @@ -44,10 +44,6 @@ import org.powermock.core.classloader.annotations.PowerMockIgnore; import org.powermock.core.classloader.annotations.PrepareForTest; import org.powermock.modules.junit4.PowerMockRunner; -import java.util.Collections; -import java.util.Enumeration; -import java.util.Locale; -import java.util.Map; import java.util.concurrent.ExecutorService; import java.util.concurrent.atomic.AtomicInteger; @@ -80,9 +76,9 @@ public class GatewayWebsocketHandlerTest { EasyMock.expect(gatewayConfig.isWebShellEnabled()).andReturn(true).anyTimes(); EasyMock.expect(gatewayConfig.getMaximumConcurrentWebshells()).andReturn(3).anyTimes(); GatewayServices gatewayServices = EasyMock.createNiceMock(GatewayServices.class); - // mock ServletUpgradeRequest and ServletUpgradeResponse - ServletUpgradeRequest req = createServletUpgradeRequest("wss://localhost:8443/gateway/webshell"); - ServletUpgradeResponse resp = createServletUpgradeResponse(); + // mock ServerUpgradeRequest and ServerUpgradeResponse + ServerUpgradeRequest req = createServerUpgradeRequest("wss://localhost:8443/gateway/webshell"); + ServerUpgradeResponse resp = createServerUpgradeResponse(); JWTValidator jwtValidator = EasyMock.createNiceMock(JWTValidator.class); EasyMock.expect(jwtValidator.validate()).andReturn(true).anyTimes(); @@ -105,9 +101,9 @@ public class GatewayWebsocketHandlerTest { EasyMock.expect(gatewayConfig.isWebShellEnabled()).andReturn(true).anyTimes(); EasyMock.expect(gatewayConfig.getMaximumConcurrentWebshells()).andReturn(3).anyTimes(); GatewayServices gatewayServices = EasyMock.createNiceMock(GatewayServices.class); - // mock ServletUpgradeRequest and ServletUpgradeResponse - ServletUpgradeRequest req = createServletUpgradeRequest("wss://www.local.com/gateway/webshell"); - ServletUpgradeResponse resp = createServletUpgradeResponse(); + // mock ServerUpgradeRequest and ServerUpgradeResponse + ServerUpgradeRequest req = createServerUpgradeRequest("wss://www.local.com/gateway/webshell"); + ServerUpgradeResponse resp = createServerUpgradeResponse(); JWTValidator jwtValidator = EasyMock.createNiceMock(JWTValidator.class); EasyMock.expect(jwtValidator.validate()).andReturn(true).anyTimes(); @@ -135,9 +131,9 @@ public class GatewayWebsocketHandlerTest { EasyMock.expect(gatewayConfig.isWebShellEnabled()).andReturn(true).anyTimes(); EasyMock.expect(gatewayConfig.getMaximumConcurrentWebshells()).andReturn(3).anyTimes(); GatewayServices gatewayServices = EasyMock.createNiceMock(GatewayServices.class); - // mock ServletUpgradeRequest and ServletUpgradeResponse - ServletUpgradeRequest req = createServletUpgradeRequest("wss://localhost:8443/gateway/webshell"); - ServletUpgradeResponse resp = createServletUpgradeResponse(); + // mock ServerUpgradeRequest and ServerUpgradeResponse + ServerUpgradeRequest req = createServerUpgradeRequest("wss://localhost:8443/gateway/webshell"); + ServerUpgradeResponse resp = createServerUpgradeResponse(); JWTValidator jwtValidator = EasyMock.createNiceMock(JWTValidator.class); EasyMock.expect(jwtValidator.validate()).andReturn(false).anyTimes(); @@ -160,79 +156,30 @@ public class GatewayWebsocketHandlerTest { GatewayConfig gatewayConfig = EasyMock.createNiceMock(GatewayConfig.class); EasyMock.expect(gatewayConfig.isWebShellEnabled()).andReturn(false).anyTimes(); GatewayServices gatewayServices = EasyMock.createNiceMock(GatewayServices.class); - // mock ServletUpgradeRequest and ServletUpgradeResponse - ServletUpgradeRequest req = createServletUpgradeRequest("wss://localhost:8443/gateway/webshell"); - ServletUpgradeResponse resp = createServletUpgradeResponse(); + // mock ServerUpgradeRequest and ServerUpgradeResponse + ServerUpgradeRequest req = createServerUpgradeRequest("wss://localhost:8443/gateway/webshell"); + ServerUpgradeResponse resp = createServerUpgradeResponse(); EasyMock.replay(gatewayServices,gatewayConfig); GatewayWebsocketHandler gatewayWebsocketHandler = new GatewayWebsocketHandler(gatewayConfig,gatewayServices); gatewayWebsocketHandler.createWebSocket(req,resp); } - private ServletUpgradeRequest createServletUpgradeRequest(String url) throws Exception { - HttpServletRequest mockRequest = new org.apache.knox.test.mock.MockHttpServletRequest() { - @Override - public StringBuffer getRequestURL() { - return new StringBuffer(url); - } - @Override - public Enumeration<String> getHeaderNames() { - return Collections.emptyEnumeration(); - } - @Override - public Map<String, String[]> getParameterMap() { - return Collections.emptyMap(); - } - @Override - public Enumeration<String> getAttributeNames() { - return Collections.emptyEnumeration(); - } - @Override - public Enumeration<Locale> getLocales() { - return Collections.emptyEnumeration(); - } - @Override - public String getRemoteAddr() { - return "127.0.0.1"; - } - @Override - public String getRemoteHost() { - return "localhost"; - } - @Override - public int getRemotePort() { - return 1234; - } - @Override - public String getLocalAddr() { - return "127.0.0.1"; - } - @Override - public String getLocalName() { - return "localhost"; - } - @Override - public int getLocalPort() { - return 8443; - } - @Override - public String getServerName() { - return "localhost"; - } - @Override - public int getServerPort() { - return 8443; - } - @Override - public String getScheme() { - return "wss"; - } - }; - return new ServletUpgradeRequest(mockRequest); + private ServerUpgradeRequest createServerUpgradeRequest(String url) throws Exception { + ServerUpgradeRequest mockRequest = EasyMock.createNiceMock(ServerUpgradeRequest.class); + HttpURI httpURI = HttpURI.build(url); + + // Set the expectations needed by KnoxWebSocketCreator and JWTValidator + EasyMock.expect(mockRequest.getHttpURI()).andReturn(httpURI).anyTimes(); + EasyMock.expect(mockRequest.getHeaders()).andReturn(HttpFields.EMPTY).anyTimes(); + + EasyMock.replay(mockRequest); + return mockRequest; } - private ServletUpgradeResponse createServletUpgradeResponse() { - HttpServletResponse mockResponse = EasyMock.createNiceMock(HttpServletResponse.class); + private ServerUpgradeResponse createServerUpgradeResponse() { + ServerUpgradeResponse mockResponse = EasyMock.createNiceMock(ServerUpgradeResponse.class); + EasyMock.replay(mockResponse); - return new ServletUpgradeResponse(mockResponse); + return mockResponse; } } diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/MessageFailureTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/MessageFailureTest.java index 3b1f94fa3..b31201ab5 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/MessageFailureTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/MessageFailureTest.java @@ -34,6 +34,7 @@ import jakarta.websocket.ContainerProvider; import jakarta.websocket.WebSocketContainer; import java.net.URI; import java.util.Locale; +import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -130,8 +131,7 @@ public class MessageFailureTest { backend.addConnector(connector); /* start backend with Echo socket */ - final BigEchoSocketHandler wsHandler = new BigEchoSocketHandler( - new EchoSocket()); + final BigEchoSocketHandler wsHandler = new BigEchoSocketHandler(EchoSocket::new); ContextHandler context = new ContextHandler(); context.setContextPath("/"); @@ -156,8 +156,10 @@ public class MessageFailureTest { proxy.addConnector(proxyConnector); /* start Knox with WebsocketAdapter to test */ + final ExecutorService pool = Executors.newFixedThreadPool(10); + final BigEchoSocketHandler wsHandler = new BigEchoSocketHandler( - new ProxyWebSocketAdapter(serverUri, Executors.newFixedThreadPool(10), gatewayConfig)); + () -> new ProxyWebSocketAdapter(serverUri, pool, gatewayConfig)); ContextHandler context = new ContextHandler(); context.setContextPath("/"); diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoHandler.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoHandler.java index 1049c9bb7..002504d7d 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoHandler.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoHandler.java @@ -17,27 +17,25 @@ */ package org.apache.knox.gateway.websockets; -import org.eclipse.jetty.websocket.server.WebSocketHandler; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeRequest; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeResponse; -import org.eclipse.jetty.websocket.servlet.WebSocketCreator; -import org.eclipse.jetty.websocket.servlet.WebSocketServletFactory; +import org.eclipse.jetty.util.Callback; +import org.eclipse.jetty.websocket.server.ServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.ServerUpgradeResponse; +import org.eclipse.jetty.websocket.server.ServerWebSocketContainer; /** - * A Mock websocket handler that just Echos messages + * A websocket handler that just Echos messages * */ -public class WebsocketEchoHandler extends WebSocketHandler implements WebSocketCreator { - private final EchoSocket socket = new EchoSocket(); +public class WebsocketEchoHandler extends AbstractWebSocketHandler { - @Override - public void configure(WebSocketServletFactory factory) { - factory.getPolicy().setMaxTextMessageSize(2 * 1024 * 1024); - factory.setCreator(this); - } + @Override + protected void configure(ServerWebSocketContainer container) { + container.setMaxTextMessageSize(2 * 1024 * 1024); + } - @Override - public Object createWebSocket(ServletUpgradeRequest req, ServletUpgradeResponse resp) { - return socket; - } + @Override + public Object createWebSocket(ServerUpgradeRequest req, ServerUpgradeResponse resp, Callback callback) { + // Must return a new instance per connection, as Session is stateful + return new EchoSocket(); + } } diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoTestBase.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoTestBase.java index 3c26466bd..9f857e0c6 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoTestBase.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketEchoTestBase.java @@ -33,11 +33,11 @@ import org.apache.knox.gateway.topology.TopologyEvent; import org.apache.knox.gateway.topology.TopologyListener; import org.apache.knox.test.TestUtils; import org.easymock.EasyMock; +import org.eclipse.jetty.server.Handler; import org.eclipse.jetty.server.Server; import org.eclipse.jetty.server.ServerConnector; import org.eclipse.jetty.server.handler.ContextHandler; -import org.eclipse.jetty.server.handler.HandlerCollection; -import org.eclipse.jetty.websocket.server.WebSocketHandler; +import org.eclipse.jetty.server.handler.ContextHandlerCollection; import java.io.File; import java.io.IOException; @@ -94,7 +94,7 @@ public class WebsocketEchoTestBase { */ public static URI serverUri; - public static WebSocketHandler handler; + public static Handler handler; private static File topoDir; private static Path dataDir; @@ -185,7 +185,7 @@ public class WebsocketEchoTestBase { gatewayServer.addConnector(connector); /* workaround so we can add our handler later at runtime */ - HandlerCollection handlers = new HandlerCollection(true); + ContextHandlerCollection handlers = new ContextHandlerCollection(true); /* add some initial handlers */ ContextHandler context = new ContextHandler(); diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketMultipleConnectionTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketMultipleConnectionTest.java index 0a7bd45f9..d104b8bb4 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketMultipleConnectionTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketMultipleConnectionTest.java @@ -43,7 +43,7 @@ import org.easymock.EasyMock; import org.eclipse.jetty.server.Server; import org.eclipse.jetty.server.ServerConnector; import org.eclipse.jetty.server.handler.ContextHandler; -import org.eclipse.jetty.server.handler.HandlerCollection; +import org.eclipse.jetty.server.handler.ContextHandlerCollection; import org.eclipse.jetty.util.thread.QueuedThreadPool; import org.junit.AfterClass; import org.junit.BeforeClass; @@ -218,7 +218,7 @@ public class WebsocketMultipleConnectionTest { gatewayServer.addConnector(connector); /* workaround so we can add our handler later at runtime */ - HandlerCollection handlers = new HandlerCollection(true); + ContextHandlerCollection handlers = new ContextHandlerCollection(true); /* add some initial handlers */ ContextHandler context = new ContextHandler(); diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedMessageTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedMessageTest.java index 5ff6b869e..938f4a45d 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedMessageTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedMessageTest.java @@ -17,23 +17,17 @@ */ package org.apache.knox.gateway.websockets; -import org.eclipse.jetty.io.RuntimeIOException; -import org.eclipse.jetty.websocket.api.BatchMode; -import org.eclipse.jetty.websocket.api.RemoteEndpoint; -import org.eclipse.jetty.websocket.api.WebSocketAdapter; +import jakarta.websocket.ContainerProvider; +import jakarta.websocket.WebSocketContainer; +import org.eclipse.jetty.util.Callback; import org.eclipse.jetty.websocket.api.Session; -import org.eclipse.jetty.websocket.server.WebSocketHandler; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeRequest; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeResponse; -import org.eclipse.jetty.websocket.servlet.WebSocketCreator; -import org.eclipse.jetty.websocket.servlet.WebSocketServletFactory; +import org.eclipse.jetty.websocket.server.ServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.ServerUpgradeResponse; +import org.eclipse.jetty.websocket.server.ServerWebSocketContainer; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.Test; -import jakarta.websocket.ContainerProvider; -import jakarta.websocket.WebSocketContainer; -import java.io.IOException; import java.net.URI; import java.util.concurrent.TimeUnit; @@ -99,25 +93,23 @@ public class WebsocketServerInitiatedMessageTest extends WebsocketEchoTestBase { * A Mock websocket handler * */ - private static class WebsocketServerInitiatedEchoHandler extends WebSocketHandler implements WebSocketCreator { - private final ServerInitiatingMessageSocket socket = new ServerInitiatingMessageSocket(); + private static class WebsocketServerInitiatedEchoHandler extends AbstractWebSocketHandler { @Override - public void configure(WebSocketServletFactory factory) { - factory.getPolicy().setMaxTextMessageSize(2 * 1024 * 1024); - factory.setCreator(this); + protected void configure(ServerWebSocketContainer container) { + container.setMaxTextMessageSize(2 * 1024 * 1024); } @Override - public Object createWebSocket(ServletUpgradeRequest req, ServletUpgradeResponse resp) { - return socket; + public Object createWebSocket(ServerUpgradeRequest req, ServerUpgradeResponse resp, Callback callback) { + return new ServerInitiatingMessageSocket(); } } /** * A simple socket initiating message on connect */ - private static class ServerInitiatingMessageSocket extends WebSocketAdapter { + private static class ServerInitiatingMessageSocket extends Session.Listener.AbstractAutoDemanding { @Override public void onWebSocketError(Throwable cause) { @@ -125,18 +117,13 @@ public class WebsocketServerInitiatedMessageTest extends WebsocketEchoTestBase { } @Override - public void onWebSocketConnect(Session sess) { - super.onWebSocketConnect(sess); - - try { - RemoteEndpoint remote = getRemote(); - remote.sendString("echo", null); - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } - } catch (IOException x) { - throw new RuntimeIOException(x); - } + public void onWebSocketOpen(Session session) { + super.onWebSocketOpen(session); + + // In Jetty 12, we send the text directly on the session and provide a Callback. + // We use Callback.NOOP since the original code passed null and ignored success/failure. + // BatchMode and manual flushing are handled automatically by the Jetty engine. + session.sendText("echo", org.eclipse.jetty.websocket.api.Callback.NOOP); } } } diff --git a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedPingTest.java b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedPingTest.java index f19b7d30d..9f38d6f57 100644 --- a/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedPingTest.java +++ b/gateway-server/src/test/java/org/apache/knox/gateway/websockets/WebsocketServerInitiatedPingTest.java @@ -17,23 +17,18 @@ */ package org.apache.knox.gateway.websockets; -import org.eclipse.jetty.io.RuntimeIOException; -import org.eclipse.jetty.websocket.api.BatchMode; -import org.eclipse.jetty.websocket.api.RemoteEndpoint; -import org.eclipse.jetty.websocket.api.WebSocketAdapter; +import org.eclipse.jetty.util.Callback; import org.eclipse.jetty.websocket.api.Session; -import org.eclipse.jetty.websocket.server.WebSocketHandler; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeRequest; -import org.eclipse.jetty.websocket.servlet.ServletUpgradeResponse; -import org.eclipse.jetty.websocket.servlet.WebSocketCreator; -import org.eclipse.jetty.websocket.servlet.WebSocketServletFactory; +import org.eclipse.jetty.websocket.server.ServerUpgradeRequest; +import org.eclipse.jetty.websocket.server.ServerUpgradeResponse; +import org.eclipse.jetty.websocket.server.ServerWebSocketContainer; + import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.Test; import jakarta.websocket.ContainerProvider; import jakarta.websocket.WebSocketContainer; -import java.io.IOException; import java.net.URI; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; @@ -101,25 +96,23 @@ public class WebsocketServerInitiatedPingTest extends WebsocketEchoTestBase { * A Mock websocket handler * */ - private static class WebsocketServerInitiatedPingHandler extends WebSocketHandler implements WebSocketCreator { - private final ServerInitiatingPingSocket socket = new ServerInitiatingPingSocket(); + private static class WebsocketServerInitiatedPingHandler extends AbstractWebSocketHandler { @Override - public void configure(WebSocketServletFactory factory) { - factory.getPolicy().setMaxTextMessageSize(2 * 1024 * 1024); - factory.setCreator(this); + protected void configure(ServerWebSocketContainer container) { + container.setMaxTextMessageSize(2 * 1024 * 1024); } @Override - public Object createWebSocket(ServletUpgradeRequest req, ServletUpgradeResponse resp) { - return socket; + public Object createWebSocket(ServerUpgradeRequest req, ServerUpgradeResponse resp, Callback callback) { + return new ServerInitiatingPingSocket(); } } /** * A simple socket initiating message on connect */ - private static class ServerInitiatingPingSocket extends WebSocketAdapter { + private static class ServerInitiatingPingSocket extends Session.Listener.AbstractAutoDemanding { @Override public void onWebSocketError(Throwable cause) { @@ -127,25 +120,23 @@ public class WebsocketServerInitiatedPingTest extends WebsocketEchoTestBase { } @Override - public void onWebSocketConnect(Session sess) { - super.onWebSocketConnect(sess); + public void onWebSocketOpen(Session session) { + super.onWebSocketOpen(session); + try { Thread.sleep(1000); - } catch (Exception e) { + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); } + final String textMessage = "PingPong"; final ByteBuffer binaryMessage = ByteBuffer.wrap( - textMessage.getBytes(StandardCharsets.UTF_8)); + textMessage.getBytes(StandardCharsets.UTF_8)); - try { - RemoteEndpoint remote = getRemote(); - remote.sendPing(binaryMessage); - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } - } catch (IOException x) { - throw new RuntimeIOException(x); - } + // In Jetty 12, we send the ping directly on the session and provide a Callback. + // We use Callback.NOOP since the original code didn't do anything on success/failure. + // BatchMode and manual flushing are handled automatically by the Jetty engine. + session.sendPing(binaryMessage, org.eclipse.jetty.websocket.api.Callback.NOOP); } } }
