From 186c2321a5bffcc65fdfe0aaf454eda39d923cad Mon Sep 17 00:00:00 2001 From: Aayush Atharva Date: Wed, 23 Sep 2026 20:50:29 +0000 Subject: [PATCH] Close with 1009 when a WebSocket message is too big --- .../netty/handler/WebSocketHandler.java | 15 ++++ .../ws/WebSocketMessageTooBigTest.java | 79 +++++++++++++++++++ 2 files changed, 94 insertions(+) create mode 100644 client/src/test/java/org/asynchttpclient/ws/WebSocketMessageTooBigTest.java diff --git a/client/src/main/java/org/asynchttpclient/netty/handler/WebSocketHandler.java b/client/src/main/java/org/asynchttpclient/netty/handler/WebSocketHandler.java index a3b7b3173..9f558ce23 100755 --- a/client/src/main/java/org/asynchttpclient/netty/handler/WebSocketHandler.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/WebSocketHandler.java @@ -17,11 +17,15 @@ import io.netty.channel.Channel; import io.netty.channel.ChannelHandler.Sharable; +import io.netty.channel.ChannelHandlerContext; +import io.netty.handler.codec.TooLongFrameException; import io.netty.handler.codec.http.HttpHeaderValues; import io.netty.handler.codec.http.HttpHeaders; import io.netty.handler.codec.http.HttpRequest; import io.netty.handler.codec.http.HttpResponse; import io.netty.handler.codec.http.LastHttpContent; +import io.netty.handler.codec.http.websocketx.CloseWebSocketFrame; +import io.netty.handler.codec.http.websocketx.WebSocketCloseStatus; import io.netty.handler.codec.http.websocketx.WebSocketFrame; import org.asynchttpclient.AsyncHandler.State; import org.asynchttpclient.AsyncHttpClientConfig; @@ -41,6 +45,8 @@ import static io.netty.handler.codec.http.HttpHeaderNames.SEC_WEBSOCKET_KEY; import static io.netty.handler.codec.http.HttpHeaderNames.UPGRADE; import static io.netty.handler.codec.http.HttpResponseStatus.SWITCHING_PROTOCOLS; +import static org.asynchttpclient.netty.channel.ChannelManager.WS_FRAME_AGGREGATOR; +import static org.asynchttpclient.util.MiscUtils.getCause; import static org.asynchttpclient.ws.WebSocketUtils.getAcceptKey; @Sharable @@ -153,6 +159,15 @@ public void handleRead(Channel channel, NettyResponseFuture future, Object e) } } + @Override + public void exceptionCaught(ChannelHandlerContext ctx, Throwable e) { + // RFC 6455 section 7.4.1: 1009 for a message too big to process. Queued before the close below. + if (getCause(e) instanceof TooLongFrameException && ctx.pipeline().get(WS_FRAME_AGGREGATOR) != null) { + ctx.writeAndFlush(new CloseWebSocketFrame(WebSocketCloseStatus.MESSAGE_TOO_BIG)); + } + super.exceptionCaught(ctx, e); + } + @Override public void handleException(NettyResponseFuture future, Throwable e) { logger.warn("onError", e); diff --git a/client/src/test/java/org/asynchttpclient/ws/WebSocketMessageTooBigTest.java b/client/src/test/java/org/asynchttpclient/ws/WebSocketMessageTooBigTest.java new file mode 100644 index 000000000..58d9fd4f6 --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/ws/WebSocketMessageTooBigTest.java @@ -0,0 +1,79 @@ +/* + * Copyright (c) 2026 AsyncHttpClient Project. All rights reserved. + * + * Licensed 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.asynchttpclient.ws; + +import org.asynchttpclient.AsyncHttpClient; +import org.eclipse.jetty.server.handler.AbstractHandler; +import org.eclipse.jetty.servlet.ServletContextHandler; +import org.eclipse.jetty.websocket.api.Session; +import org.eclipse.jetty.websocket.api.WebSocketAdapter; +import org.eclipse.jetty.websocket.server.config.JettyWebSocketServletContainerInitializer; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +import java.io.IOException; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; + +import static org.asynchttpclient.Dsl.asyncHttpClient; +import static org.asynchttpclient.Dsl.config; +import static org.junit.jupiter.api.Assertions.assertEquals; + +public class WebSocketMessageTooBigTest extends AbstractBasicWebSocketTest { + + private static volatile CompletableFuture serverCloseCode = new CompletableFuture<>(); + + public static class FragmentedSender extends WebSocketAdapter { + + @Override + public void onWebSocketConnect(Session session) { + super.onWebSocketConnect(session); + try { + String fragment = "x".repeat(1000); + for (int i = 0; i < 4; i++) { + getRemote().sendPartialString(fragment, i == 3); + } + } catch (IOException e) { + serverCloseCode.completeExceptionally(e); + } + } + + @Override + public void onWebSocketClose(int statusCode, String reason) { + serverCloseCode.complete(statusCode); + super.onWebSocketClose(statusCode, reason); + } + } + + @Override + public AbstractHandler configureHandler() { + ServletContextHandler context = new ServletContextHandler(ServletContextHandler.SESSIONS); + context.setContextPath("/"); + JettyWebSocketServletContainerInitializer.configure(context, + (servletContext, wsContainer) -> wsContainer.addMapping("/", FragmentedSender.class)); + return context; + } + + @Test + @Timeout(unit = TimeUnit.MILLISECONDS, value = 30000) + public void aMessageOverTheBufferLimitClosesWith1009() throws Exception { + serverCloseCode = new CompletableFuture<>(); + try (AsyncHttpClient client = asyncHttpClient(config().setWebSocketMaxBufferSize(1024))) { + client.prepareGet(getTargetUrl()).execute(new WebSocketUpgradeHandler.Builder().build()).get(); + assertEquals(1009, serverCloseCode.get(10, TimeUnit.SECONDS)); + } + } +}