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 b53fcd43de7943a5f70abd8978defd0db792eb33 Author: bonampak <[email protected]> AuthorDate: Fri Apr 24 03:49:04 2026 +0200 KNOX-3238: correct GatewayWebsocketHandler, ProxyWebSocketAdapter, WebshellWebSocketAdapter and JWTValidatorFactory. --- gateway-server/pom.xml | 5 + .../org/apache/knox/gateway/GatewayServer.java | 2 +- .../gateway/webshell/WebshellWebSocketAdapter.java | 35 ++- .../websockets/GatewayWebsocketHandler.java | 214 +++++++++++------ .../gateway/websockets/JWTValidatorFactory.java | 12 +- .../gateway/websockets/ProxyWebSocketAdapter.java | 262 ++++++++++----------- 6 files changed, 310 insertions(+), 220 deletions(-) diff --git a/gateway-server/pom.xml b/gateway-server/pom.xml index 5a4104da2..c67b0ae25 100644 --- a/gateway-server/pom.xml +++ b/gateway-server/pom.xml @@ -339,6 +339,11 @@ <groupId>org.eclipse.jetty.websocket</groupId> <artifactId>jetty-websocket-jetty-api</artifactId> </dependency> + <dependency> + <groupId>org.eclipse.jetty.websocket</groupId> + <artifactId>jetty-websocket-jetty-server</artifactId> + </dependency> + <dependency> <groupId>org.eclipse.jetty.ee10.websocket</groupId> <artifactId>jetty-ee10-websocket-jetty-server</artifactId> diff --git a/gateway-server/src/main/java/org/apache/knox/gateway/GatewayServer.java b/gateway-server/src/main/java/org/apache/knox/gateway/GatewayServer.java index 85240f2bf..5fe803141 100644 --- a/gateway-server/src/main/java/org/apache/knox/gateway/GatewayServer.java +++ b/gateway-server/src/main/java/org/apache/knox/gateway/GatewayServer.java @@ -1293,7 +1293,7 @@ public class GatewayServer { request.setAttribute(ErrorHandler.ERROR_MESSAGE, "Service Unavailable"); // The default ErrorHandler references this attribute response.setStatus(HttpServletResponse.SC_SERVICE_UNAVAILABLE); } - super.handle(request, response, callback); + return super.handle(request, response, callback); } } diff --git a/gateway-server/src/main/java/org/apache/knox/gateway/webshell/WebshellWebSocketAdapter.java b/gateway-server/src/main/java/org/apache/knox/gateway/webshell/WebshellWebSocketAdapter.java index 44fb8d1a6..7ec9e4ac9 100644 --- a/gateway-server/src/main/java/org/apache/knox/gateway/webshell/WebshellWebSocketAdapter.java +++ b/gateway-server/src/main/java/org/apache/knox/gateway/webshell/WebshellWebSocketAdapter.java @@ -18,6 +18,7 @@ package org.apache.knox.gateway.webshell; import java.io.IOException; +import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.util.concurrent.ExecutorService; import java.util.concurrent.atomic.AtomicInteger; @@ -35,7 +36,9 @@ import org.apache.knox.gateway.config.GatewayConfig; import org.apache.knox.gateway.services.security.token.UnknownTokenException; import org.apache.knox.gateway.websockets.JWTValidator; import org.apache.knox.gateway.websockets.ProxyWebSocketAdapter; +import org.eclipse.jetty.websocket.api.Callback; import org.eclipse.jetty.websocket.api.Session; +import org.eclipse.jetty.websocket.api.StatusCode; public class WebshellWebSocketAdapter extends ProxyWebSocketAdapter { private Session session; @@ -61,7 +64,7 @@ public class WebshellWebSocketAdapter extends ProxyWebSocketAdapter { @SuppressWarnings("PMD.DoNotUseThreads") @Override - public void onWebSocketConnect(final Session session) { + public void onWebSocketOpen(final Session session) { this.session = session; connectionInfo.connect(); pool.execute(this::blockingReadFromHost); @@ -111,18 +114,28 @@ public class WebshellWebSocketAdapter extends ProxyWebSocketAdapter { } private void transToClient(String message){ - try { - session.getRemote().sendString(message); - } catch (IOException e){ - LOG.onError("Error sending message to client"); + if (session == null) { + LOG.onError("Cannot send message to client; session is null"); cleanup(); + return; } + // In Jetty 12 the send is asynchronous: report send failures via the Callback, + // preserving the old IOException-based error handling semantics. + session.sendText(message, Callback.from( + () -> { /* success: nothing to do */ }, + t -> { + LOG.onError("Error sending message to client"); + cleanup(); + })); } @Override - public void onWebSocketBinary(final byte[] payload, final int offset, final int length) { - throw new UnsupportedOperationException( - "Websocket for binary messages is not supported at this time."); + public void onWebSocketBinary(final ByteBuffer payload, final Callback callback) { + // Binary is not supported for webshell sessions. Complete the callback so the + // Jetty 12 demand loop can surface the error and close the connection cleanly + // rather than throwing from the listener (which would leave demand stalled). + callback.fail(new UnsupportedOperationException( + "Websocket for binary messages is not supported at this time.")); } @Override @@ -166,10 +179,10 @@ public class WebshellWebSocketAdapter extends ProxyWebSocketAdapter { } private void cleanup() { - if(session != null && session.isOpen()) { - session.close(); + if (session != null && session.isOpen()) { + session.close(StatusCode.NORMAL, null, Callback.NOOP); session = null; } connectionInfo.disconnect(); } -} +} \ No newline at end of file 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 a52f42551..3e3d0ce85 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 @@ -28,18 +28,30 @@ import org.apache.knox.gateway.services.registry.ServiceRegistry; import org.apache.knox.gateway.services.security.KeystoreService; import org.apache.knox.gateway.services.security.KeystoreServiceException; import org.apache.knox.gateway.webshell.WebshellWebSocketAdapter; -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; + +// Jetty 12 Core & WebSocket Imports +import org.eclipse.jetty.http.HttpField; +import org.eclipse.jetty.http.HttpURI; +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.WebSocketCreator; +import org.eclipse.jetty.websocket.server.ServerWebSocketContainer; +import org.eclipse.jetty.websocket.server.WebSocketUpgradeHandler; + import jakarta.websocket.ClientEndpointConfig; import java.net.MalformedURLException; import java.net.URI; import java.net.URISyntaxException; import java.net.URL; import java.security.KeyStore; -import java.util.Arrays; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.Locale; import java.util.Map; @@ -54,8 +66,7 @@ import java.util.concurrent.atomic.AtomicInteger; * * @since 0.10 */ -public class GatewayWebsocketHandler extends WebSocketHandler - implements WebSocketCreator { +public class GatewayWebsocketHandler extends Handler.Wrapper { private static final WebsocketLogMessages LOG = MessagesFactory .get(WebsocketLogMessages.class); @@ -82,7 +93,7 @@ public class GatewayWebsocketHandler extends WebSocketHandler final GatewayConfig config; final GatewayServices services; - + private WebSocketUpgradeHandler wsHandler; public GatewayWebsocketHandler(final GatewayConfig config, final GatewayServices services) { @@ -91,27 +102,108 @@ public class GatewayWebsocketHandler extends WebSocketHandler this.services = services; pool = Executors.newFixedThreadPool(POOL_SIZE); this.concurrentWebshells = new AtomicInteger(0); + // Set the internal handler as the one we are wrapping + setHandler(wsHandler); + } + + @Override + protected void doStart() throws Exception { + Server server = getServer(); + if (server == null) { + throw new IllegalStateException("GatewayWebsocketHandler must be attached to a Server before starting"); + } + + // 1. Get or create the global ServerWebSocketContainer (no ContextHandler needed) + ServerWebSocketContainer container = ServerWebSocketContainer.ensure(server, null); + configureServerWebSocketContainer(container); + + // 3. Create the UpgradeHandler. + this.wsHandler = new WebSocketUpgradeHandler(container); + + // 4. PRESERVE THE CHAIN: This handler currently wraps PortMappingHelperHandler. + // We must ensure the internal wsHandler also wraps it, so HTTP traffic flows downwards. + this.wsHandler.setHandler(getHandler()); + + // Start the internal handler + this.wsHandler.start(); + + super.doStart(); + } + + @Override + protected void doStop() throws Exception { + if (this.wsHandler != null) { + this.wsHandler.stop(); + } + super.doStop(); } @Override - public void configure(final WebSocketServletFactory factory) { - factory.setCreator(this); - factory.getPolicy() - .setMaxTextMessageSize(config.getWebsocketMaxTextMessageSize()); - factory.getPolicy() - .setMaxBinaryMessageSize(config.getWebsocketMaxBinaryMessageSize()); + public boolean handle(Request request, Response response, Callback callback) throws Exception { + // Delegate to the WebSocketUpgradeHandler. + // If it detects a WebSocket upgrade, it triggers KnoxWebSocketCreator and returns true. + // If it's standard HTTP traffic, it delegates to its wrapped handler (PortMappingHelperHandler). + return wsHandler.handle(request, response, callback); + } - factory.getPolicy().setMaxBinaryMessageBufferSize( - config.getWebsocketMaxBinaryMessageBufferSize()); - factory.getPolicy().setMaxTextMessageBufferSize( - config.getWebsocketMaxTextMessageBufferSize()); + private class KnoxWebSocketCreator implements WebSocketCreator { + @Override + public Object createWebSocket(ServerUpgradeRequest req, ServerUpgradeResponse resp, Callback callback) { + try { + // 1. Get the raw HTTP URI from the Jetty 12 Request + HttpURI httpURI = req.getHttpURI(); - factory.getPolicy() - .setInputBufferSize(config.getWebsocketInputBufferSize()); + // 2. Translate the scheme to match Jetty 9's behavior (http -> ws, https -> wss) + String wsScheme = "https".equalsIgnoreCase(httpURI.getScheme()) ? "wss" : "ws"; + + // 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! + if (isWebshellRequest(requestURI)) { + return handleWebshellRequest(req); // Note: Update handleWebshellRequest to accept ServerUpgradeRequest + } - factory.getPolicy() - .setAsyncWriteTimeout(config.getWebsocketAsyncWriteTimeout()); - factory.getPolicy().setIdleTimeout(config.getWebsocketIdleTimeout()); + final String backendURL = getMatchedBackendURL(requestURI); + LOG.debugLog("Generated backend URL for websocket connection: " + backendURL); + + final ClientEndpointConfig clientConfig = getClientEndpointConfig(req, backendURL); + clientConfig.getUserProperties().put("org.apache.knox.gateway.websockets.truststore", getTruststore()); + + return new ProxyWebSocketAdapter(URI.create(backendURL), pool, clientConfig, config); + + } catch (final Exception e) { + LOG.failedCreatingWebSocket(e); + // In Jetty 12, completing the callback with failure tells the server to reject the upgrade + callback.failed(e); + return null; + } + } + } + + public void configureServerWebSocketContainer(ServerWebSocketContainer container) { + container.setMaxTextMessageSize(config.getWebsocketMaxTextMessageSize()); + container.setMaxBinaryMessageSize(config.getWebsocketMaxBinaryMessageSize()); + container.setInputBufferSize(config.getWebsocketInputBufferSize()); + + container.setIdleTimeout(Duration.ofMillis(config.getWebsocketIdleTimeout())); + + // 2. Map ALL incoming requests to our custom Knox routing creator + // "regex|^/.*" acts as a catch-all interceptor. + container.addMapping("regex|^/.*", new KnoxWebSocketCreator()); + + //removed in Jetty 12 container.setMaxBinaryMessageBufferSize(config.getWebsocketMaxBinaryMessageBufferSize()); + //removed in Jetty 12 container.setMaxTextMessageBufferSize(config.getWebsocketMaxTextMessageBufferSize()); + + //removed in Jetty 12 container.setAsyncWriteTimeout(config.getWebsocketAsyncWriteTimeout()); + // handled by the core HTTP connection idle timeouts and + // one can apply it directly to the asynchronous write execution + // (e.g., using CompletableFuture.orTimeout(duration, TimeUnit) when we call session.sendText(...)). + + // removed, idle timeout is used or specified in send() methods: + // container.setAsyncSendTimeout(config.getWebsocketAsyncWriteTimeout()); + + //same as setIdleTimeout: container.setDefaultMaxSessionIdleTimeout(config.getWebsocketIdleTimeout()); } @@ -119,7 +211,7 @@ public class GatewayWebsocketHandler extends WebSocketHandler return requestURI.toString().matches(REGEX_WEBSHELL_REQUEST_PATH); } - private WebshellWebSocketAdapter handleWebshellRequest(ServletUpgradeRequest req){ + private WebshellWebSocketAdapter handleWebshellRequest(ServerUpgradeRequest req){ if (config.isWebShellEnabled()){ if (concurrentWebshells.get() >= config.getMaximumConcurrentWebshells()){ throw new RuntimeException("Number of allowed concurrent Web Shell sessions exceeded"); @@ -133,31 +225,6 @@ public class GatewayWebsocketHandler extends WebSocketHandler throw new RuntimeException("Web Shell not enabled"); } - - @Override - public Object createWebSocket(ServletUpgradeRequest req, - ServletUpgradeResponse resp) { - try { - final URI requestURI = req.getRequestURI(); - - if (isWebshellRequest(requestURI)) { - return handleWebshellRequest(req); - } - - // URL used to connect to websocket backend - final String backendURL = getMatchedBackendURL(requestURI); - LOG.debugLog("Generated backend URL for websocket connection: " + backendURL); - - // Upgrade happens here - final ClientEndpointConfig clientConfig = getClientEndpointConfig(req); - clientConfig.getUserProperties().put("org.apache.knox.gateway.websockets.truststore", getTruststore()); - return new ProxyWebSocketAdapter(URI.create(backendURL), pool, clientConfig, config); - } catch (final Exception e) { - LOG.failedCreatingWebSocket(e); - throw new RuntimeException(e); - } - } - private KeyStore getTruststore() throws KeystoreServiceException { final KeystoreService ks = this.services .getService(ServiceType.KEYSTORE_SERVICE); @@ -175,26 +242,37 @@ public class GatewayWebsocketHandler extends WebSocketHandler * to be passed to the backend. * @since 0.14.0 */ - private ClientEndpointConfig getClientEndpointConfig(final ServletUpgradeRequest req) { + private ClientEndpointConfig getClientEndpointConfig(final ServerUpgradeRequest req, final String backendURL) { return ClientEndpointConfig.Builder.create() - .configurator(new ClientEndpointConfig.Configurator() { - - @Override - public void beforeRequest(final Map<String, List<String>> headers) { - - /* Add request headers */ - req.getHeaders().forEach(headers::putIfAbsent); - try { - final URI backendURL = new URI(getMatchedBackendURL(req.getRequestURI())); - headers.put("Host", Arrays.asList(backendURL.getHost() + ":" + backendURL.getPort())); - } catch (final URISyntaxException e) { - LOG.onError(String.format(Locale.ROOT, - "Error getting backend url, this could cause 'Host does not match SNI' exception. Cause: ", - e.toString())); - } - } - }).build(); + .configurator(new ClientEndpointConfig.Configurator() { + + @Override + public void beforeRequest(final Map<String, List<String>> headers) { + + // 1. Safely iterate over Jetty 12 HttpFields and copy them to the Jakarta map + for (HttpField field : req.getHeaders()) { + headers.computeIfAbsent(field.getName(), k -> new ArrayList<>()) + .add(field.getValue()); + } + + // 2. Properly construct and override the Host header + try { + final URI backendURI = new URI(backendURL); + + // Handle implicit ports (where getPort() returns -1) to prevent "Host: example.com:-1" + int port = backendURI.getPort(); + String hostValue = backendURI.getHost() + (port != -1 ? ":" + port : ""); + + headers.put("Host", Collections.singletonList(hostValue)); + + } catch (final URISyntaxException e) { + LOG.onError(String.format(Locale.ROOT, + "Error getting backend url, this could cause 'Host does not match SNI' exception. Cause: %s", + e.toString())); + } + } + }).build(); } /** diff --git a/gateway-server/src/main/java/org/apache/knox/gateway/websockets/JWTValidatorFactory.java b/gateway-server/src/main/java/org/apache/knox/gateway/websockets/JWTValidatorFactory.java index bf9740882..658a2ae91 100644 --- a/gateway-server/src/main/java/org/apache/knox/gateway/websockets/JWTValidatorFactory.java +++ b/gateway-server/src/main/java/org/apache/knox/gateway/websockets/JWTValidatorFactory.java @@ -30,9 +30,11 @@ import org.apache.knox.gateway.services.topology.TopologyService; import org.apache.knox.gateway.topology.Service; import org.apache.knox.gateway.topology.Topology; import org.apache.knox.gateway.util.CertificateUtils; -import org.eclipse.jetty.ee10.websocket.server.JettyServerUpgradeRequest; import jakarta.servlet.ServletException; -import java.net.HttpCookie; +import org.eclipse.jetty.http.HttpCookie; +import org.eclipse.jetty.server.Request; +import org.eclipse.jetty.websocket.server.ServerUpgradeRequest; + import java.security.interfaces.RSAPublicKey; import java.text.ParseException; import java.util.LinkedHashMap; @@ -47,7 +49,7 @@ public class JWTValidatorFactory { public static final String SSO_VERIFICATION_PEM = "sso.token.verification.pem"; private static final JWTMessages jwtMessagesLog = MessagesFactory.get(JWTMessages.class); - public static JWTValidator create(JettyServerUpgradeRequest req, GatewayServices gatewayServices, + public static JWTValidator create(ServerUpgradeRequest req, GatewayServices gatewayServices, GatewayConfig gatewayConfig){ Map<String,String> params = getParams(gatewayServices); String cookieName = params.containsKey(KNOXSSO_COOKIE_NAME)? params.get(KNOXSSO_COOKIE_NAME):DEFAULT_SSO_COOKIE_NAME; @@ -97,8 +99,8 @@ public class JWTValidatorFactory { return params; } - private static JWT extractToken(JettyServerUpgradeRequest req, String cookieName){ - List<HttpCookie> ssoCookies = req.getCookies(); + private static JWT extractToken(ServerUpgradeRequest req, String cookieName){ + List<HttpCookie> ssoCookies = Request.getCookies(req); if (ssoCookies != null){ for (HttpCookie ssoCookie : ssoCookies) { if (cookieName.equals(ssoCookie.getName())) { diff --git a/gateway-server/src/main/java/org/apache/knox/gateway/websockets/ProxyWebSocketAdapter.java b/gateway-server/src/main/java/org/apache/knox/gateway/websockets/ProxyWebSocketAdapter.java index 9cc9847c9..039d454fc 100644 --- a/gateway-server/src/main/java/org/apache/knox/gateway/websockets/ProxyWebSocketAdapter.java +++ b/gateway-server/src/main/java/org/apache/knox/gateway/websockets/ProxyWebSocketAdapter.java @@ -18,54 +18,68 @@ package org.apache.knox.gateway.websockets; import java.io.IOException; +import java.io.UncheckedIOException; import java.net.URI; -import java.util.List; +import java.nio.ByteBuffer; +import java.security.KeyStore; import java.util.ArrayList; +import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import jakarta.websocket.ClientEndpointConfig; import jakarta.websocket.CloseReason; -import jakarta.websocket.ContainerProvider; import jakarta.websocket.DeploymentException; import jakarta.websocket.WebSocketContainer; -import org.apache.knox.gateway.i18n.messages.MessagesFactory; import org.apache.knox.gateway.config.GatewayConfig; -import org.eclipse.jetty.io.RuntimeIOException; +import org.apache.knox.gateway.i18n.messages.MessagesFactory; +import org.eclipse.jetty.client.HttpClient; +import org.eclipse.jetty.ee10.websocket.jakarta.client.JakartaWebSocketClientContainerProvider; +import org.eclipse.jetty.websocket.api.Callback; import org.eclipse.jetty.util.component.LifeCycle; -import org.eclipse.jetty.websocket.api.BatchMode; -import org.eclipse.jetty.websocket.api.RemoteEndpoint; +import org.eclipse.jetty.util.ssl.SslContextFactory; import org.eclipse.jetty.websocket.api.Session; import org.eclipse.jetty.websocket.api.StatusCode; -import org.eclipse.jetty.websocket.api.WebSocketAdapter; -import java.security.KeyStore; + /** * Handles outbound/inbound Websocket connections and sessions. * + * <p>The frontend (browser <-> Knox) side uses the native Jetty 12 + * WebSocket API via {@link Session.Listener.AbstractAutoDemanding}. The + * backend (Knox <-> upstream) side stays on JSR-356 ({@code jakarta.websocket}). + * * @since 0.10 */ -public class ProxyWebSocketAdapter extends WebSocketAdapter { +public class ProxyWebSocketAdapter extends Session.Listener.AbstractAutoDemanding { protected static final WebsocketLogMessages LOG = MessagesFactory.get(WebsocketLogMessages.class); - /* URI for the backend */ + private static final String TRUSTSTORE_USER_PROPERTY = + "org.apache.knox.gateway.websockets.truststore"; + + /** URI for the backend */ private final URI backend; - /* Session between the frontend (browser) and Knox */ + /** Session between the frontend (browser) and Knox */ private Session frontendSession; - /* Session between the backend (outbound) and Knox */ + /** Session between the backend (outbound) and Knox */ private jakarta.websocket.Session backendSession; + /** JSR-356 client container used to connect to the backend */ private WebSocketContainer container; protected ExecutorService pool; - /* Message buffer for holding data frames temporarily in memory till connection is setup. - Keeping the max size of the buffer as 100 messages for now. */ - private List<String> messageBuffer = new ArrayList<>(); - private Lock remoteLock = new ReentrantLock(); + /** + * Buffer for messages arriving from the backend before the frontend session + * is ready. Capped by {@code config.getWebsocketMaxWaitBufferCount()}. + */ + private final List<String> messageBuffer = new ArrayList<>(); + + /** Guards the buffer-vs-send transition at open time. */ + private final Lock remoteLock = new ReentrantLock(); protected final GatewayConfig config; @@ -73,14 +87,16 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { * Used to transmit headers from browser to backend server. * @since 0.14 */ - private ClientEndpointConfig clientConfig; + private final ClientEndpointConfig clientConfig; public ProxyWebSocketAdapter(final URI backend, final ExecutorService pool, GatewayConfig config) { this(backend, pool, null, config); } - public ProxyWebSocketAdapter(final URI backend, final ExecutorService pool, final ClientEndpointConfig clientConfig, - GatewayConfig config) { + public ProxyWebSocketAdapter(final URI backend, + final ExecutorService pool, + final ClientEndpointConfig clientConfig, + final GatewayConfig config) { super(); this.backend = backend; this.pool = pool; @@ -89,90 +105,80 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { } @Override - public void onWebSocketConnect(final Session frontEndSession) { + public void onWebSocketOpen(final Session frontEndSession) { /* - * Let's connect to the backend, this is where the Backend-to-frontend - * plumbing takes place + * Let's connect to the backend. This is where the Backend-to-Frontend + * plumbing takes place. */ - container = ContainerProvider.getWebSocketContainer(); - container.setDefaultMaxTextMessageBufferSize(frontEndSession.getPolicy().getMaxTextMessageBufferSize()); - container.setDefaultMaxBinaryMessageBufferSize(frontEndSession.getPolicy().getMaxBinaryMessageBufferSize()); - container.setAsyncSendTimeout(frontEndSession.getPolicy().getAsyncWriteTimeout()); - container.setDefaultMaxSessionIdleTimeout(frontEndSession.getPolicy().getIdleTimeout()); - - KeyStore ks = null; - if(clientConfig != null) { - ks = (KeyStore) clientConfig.getUserProperties().get("org.apache.knox.gateway.websockets.truststore"); + KeyStore truststore = null; + if (clientConfig != null) { + truststore = (KeyStore) clientConfig.getUserProperties().get(TRUSTSTORE_USER_PROPERTY); } + container = buildBackendContainer(truststore); + /* - Currently javax.websocket API has no provisions to configure SSL - https://github.com/eclipse-ee4j/websocket-api/issues/210 - Until that gets fixed we'll have to resort to this. - */ - if(container instanceof org.eclipse.jetty.websocket.jsr356.ClientContainer && - ((org.eclipse.jetty.websocket.jsr356.ClientContainer)container).getClient() != null && - ((org.eclipse.jetty.websocket.jsr356.ClientContainer)container).getClient().getSslContextFactory() != null ) { - ((org.eclipse.jetty.websocket.jsr356.ClientContainer)container).getClient().getHttpClient().getSslContextFactory().setTrustStore(ks); - LOG.logMessage("Truststore for websocket setup"); - } + * Seed sensible defaults on the outbound (JSR-356) container from the + * inbound session. In Jetty 12 the payload size / idle timeout knobs live + * on the server container (not the session) and are propagated here so + * that the backend connection honors the same limits. + */ + container.setDefaultMaxTextMessageBufferSize(config.getWebsocketMaxTextMessageBufferSize()); + container.setDefaultMaxBinaryMessageBufferSize(config.getWebsocketMaxBinaryMessageBufferSize()); + container.setAsyncSendTimeout(config.getWebsocketAsyncWriteTimeout()); + container.setDefaultMaxSessionIdleTimeout(config.getWebsocketIdleTimeout()); final ProxyInboundClient backendSocket = new ProxyInboundClient(getMessageCallback()); - /* build the configuration */ - /* Attempt Connect */ try { backendSession = container.connectToServer(backendSocket, clientConfig, backend); - LOG.onConnectionOpen(backend.toString()); - } catch (DeploymentException e) { LOG.connectionFailed(e); throw new RuntimeException(e); } catch (IOException e) { LOG.connectionFailed(e); - throw new RuntimeIOException(e); + throw new UncheckedIOException(e); } remoteLock.lock(); - super.onWebSocketConnect(frontEndSession); - this.frontendSession = frontEndSession; - - final RemoteEndpoint remote = frontEndSession.getRemote(); try { + this.frontendSession = frontEndSession; if (!messageBuffer.isEmpty()) { - flushBufferedMessages(remote); - - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } + flushBufferedMessages(); } else { LOG.debugLog("Message buffer is empty"); } - } catch (IOException e) { - LOG.connectionFailed(e); - throw new RuntimeIOException(e); - } - finally - { + } finally { remoteLock.unlock(); } } - @Override - public void onWebSocketBinary(final byte[] payload, final int offset, final int length) { - if (isNotConnected()) { - return; + /** + * Build the JSR-356 client container and, if a truststore was supplied, + * back it with an {@link HttpClient} whose SSL context trusts those certs. + * + * <p>In Jetty 12 the old {@code org.eclipse.jetty.websocket.jsr356.ClientContainer} + * cast is no longer available; the supported way to configure SSL for the + * JSR-356 client is to hand {@link JakartaWebSocketClientContainerProvider} + * a pre-configured {@link HttpClient}. + */ + private WebSocketContainer buildBackendContainer(final KeyStore truststore) { + if (truststore == null) { + return JakartaWebSocketClientContainerProvider.getContainer(null); } - - throw new UnsupportedOperationException( - "Websocket support for binary messages is not supported at this time."); + SslContextFactory.Client sslContextFactory = new SslContextFactory.Client(); + sslContextFactory.setTrustStore(truststore); + HttpClient httpClient = new HttpClient(); + httpClient.setSslContextFactory(sslContextFactory); + LOG.logMessage("Truststore for websocket setup"); + return JakartaWebSocketClientContainerProvider.getContainer(httpClient); } @Override public void onWebSocketText(final String message) { - if (isNotConnected()) { + if (frontendSession == null || !frontendSession.isOpen()) { return; } @@ -181,12 +187,21 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { /* Proxy message to backend */ try { backendSession.getBasicRemote().sendText(message); - } catch (IOException e) { LOG.connectionFailed(e); } } + @Override + public void onWebSocketBinary(final ByteBuffer payload, final Callback callback) { + if (frontendSession == null || !frontendSession.isOpen()) { + callback.succeed(); + return; + } + callback.fail(new UnsupportedOperationException( + "Websocket support for binary messages is not supported at this time.")); + } + @Override public void onWebSocketClose(int statusCode, String reason) { super.onWebSocketClose(statusCode, reason); @@ -200,20 +215,17 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { } /** - * Cleanup sessions + * Cleanup sessions on error. */ private void cleanupOnError(final Throwable t) { - LOG.onError(t.toString()); if (t.toString().contains("exceeds maximum size")) { - if(frontendSession != null && frontendSession.isOpen()) { - frontendSession.close(StatusCode.MESSAGE_TOO_LARGE, t.getMessage()); + if (frontendSession != null && frontendSession.isOpen()) { + frontendSession.close(StatusCode.MESSAGE_TOO_LARGE, t.getMessage(), Callback.NOOP); } - } - - else { - if(frontendSession != null && frontendSession.isOpen()) { - frontendSession.close(StatusCode.SERVER_ERROR, t.getMessage()); + } else { + if (frontendSession != null && frontendSession.isOpen()) { + frontendSession.close(StatusCode.SERVER_ERROR, t.getMessage(), Callback.NOOP); } cleanup(); } @@ -235,8 +247,10 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { @Override public void onConnectionClose(final CloseReason reason) { try { - frontendSession.close(reason.getCloseCode().getCode(), - reason.getReasonPhrase()); + if (frontendSession != null && frontendSession.isOpen()) { + frontendSession.close(reason.getCloseCode().getCode(), + reason.getReasonPhrase(), Callback.NOOP); + } } finally { cleanup(); } @@ -251,12 +265,12 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { public void onMessageText(String message, Object session) { LOG.logMessage("[From Backend <---]" + message); remoteLock.lock(); - final RemoteEndpoint remote = getRemote(); try { - if (remote == null) { - LOG.debugLog("Remote endpoint is null"); + if (frontendSession == null) { + LOG.debugLog("Frontend session is null"); if (messageBuffer.size() >= config.getWebsocketMaxWaitBufferCount()) { - throw new RuntimeIOException("Remote is null and message buffer is full. Cannot buffer anymore "); + throw new UncheckedIOException(new IOException( + "Frontend session is not ready and message buffer is full.")); } LOG.debugLog("Buffering message: " + message); messageBuffer.add(message); @@ -264,78 +278,52 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { } /* Proxy message to frontend */ - flushBufferedMessages(remote); + flushBufferedMessages(); LOG.debugLog("Sending current message [From Backend <---]: " + message); - remote.sendString(message); - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } - } catch (IOException e) { - LOG.connectionFailed(e); - throw new RuntimeIOException(e); - } - finally - { + frontendSession.sendText(message, Callback.NOOP); + } finally { remoteLock.unlock(); } } @Override - public void onMessageBinary(byte[] message, boolean last, - Object session) { + public void onMessageBinary(byte[] message, boolean last, Object session) { throw new UnsupportedOperationException( - "Websocket support for binary messages is not supported at this time."); - + "Websocket support for binary messages is not supported at this time."); } @Override public void onMessagePong(jakarta.websocket.PongMessage message, Object session) { - LOG.logMessage("[From Backend <---]: PING"); + LOG.logMessage("[From Backend <---]: PONG"); remoteLock.lock(); - final RemoteEndpoint remote = getRemote(); try { - if (remote == null) { - LOG.debugLog("Remote endpoint is null"); + if (frontendSession == null) { + LOG.debugLog("Frontend session is null"); return; } - /* Proxy Ping message to frontend */ - flushBufferedMessages(remote); + /* Proxy Pong message to frontend */ + flushBufferedMessages(); LOG.logMessage("Sending current PING [From Backend <---]: "); - remote.sendPing(message.getApplicationData()); - if (remote.getBatchMode() == BatchMode.ON) { - remote.flush(); - } - } catch (IOException e) { - LOG.connectionFailed(e); - throw new RuntimeIOException(e); - } - finally - { + frontendSession.sendPing(message.getApplicationData(), Callback.NOOP); + } finally { remoteLock.unlock(); } } - }; - } @SuppressWarnings("PMD.DoNotUseThreads") private void cleanup() { - /* do the cleaning business in separate thread so we don't block */ - pool.execute(new Runnable() { - @Override - public void run() { - closeQuietly(); - } - }); + /* do the cleaning business in a separate thread so we don't block */ + pool.execute(this::closeQuietly); } private void closeQuietly() { try { - if(backendSession != null && !backendSession.isOpen()) { + if (backendSession != null && backendSession.isOpen()) { backendSession.close(); } } catch (IOException e) { @@ -350,20 +338,24 @@ public class ProxyWebSocketAdapter extends WebSocketAdapter { } } - if(frontendSession != null && !frontendSession.isOpen()) { - frontendSession.close(); + if (frontendSession != null && frontendSession.isOpen()) { + frontendSession.close(StatusCode.NORMAL, null, Callback.NOOP); } } - /* - * Function to flush buffered messages. Should be called with remoteLock held + /** + * Flush buffered messages to the frontend session. Must be called while + * holding {@link #remoteLock}. */ - private void flushBufferedMessages(final RemoteEndpoint remote) throws IOException { + private void flushBufferedMessages() { + if (messageBuffer.isEmpty()) { + return; + } LOG.debugLog("Flushing old buffered messages"); - for(String obj:messageBuffer) { - LOG.debugLog("Sending old buffered message [From Backend <---]: " + obj); - remote.sendString(obj); + for (String buffered : messageBuffer) { + LOG.debugLog("Sending old buffered message [From Backend <---]: " + buffered); + frontendSession.sendText(buffered, Callback.NOOP); } messageBuffer.clear(); } -} +} \ No newline at end of file
