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 &lt;-&gt; Knox) side uses the native Jetty 12
+ * WebSocket API via {@link Session.Listener.AbstractAutoDemanding}. The
+ * backend (Knox &lt;-&gt; 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

Reply via email to