This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch backport/CAMEL-25375-camel-4.22.x in repository https://gitbox.apache.org/repos/asf/camel.git
commit 690640f1f6c2f4ba20fd3dd675c69863e0940b92 Author: Andrea Cosentino <[email protected]> AuthorDate: Tue Oct 6 14:28:54 2026 +0200 CAMEL-25375: camel-undertow - Apply securityProvider, allowedRoles and handlers to WebSocket endpoints (#27435) WebSocket (ws/wss) endpoints of camel-undertow now apply securityProvider, allowedRoles, handlers and accessLog, as the HTTP path does, and the checks fail closed. A connection must have passed the applicable consumer's handler chain to be served, and endpoint-configured providers that need a servlet context get one. Undeploying only when the server stops also fixes an NPE with two endpoints on one port. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Signed-off-by: Andrea Cosentino <[email protected]> (cherry picked from commit 3a75bfbf1b202b0f62795bdb67994bc522a807f6) --- .../camel/catalog/docs/undertow-component.adoc | 22 +- .../security/SpringSecurityWebSocketTest.java | 90 +++++++ .../src/main/docs/undertow-component.adoc | 22 +- .../component/undertow/DefaultUndertowHost.java | 97 ++++--- .../component/undertow/UndertowComponent.java | 24 +- .../camel/component/undertow/UndertowConsumer.java | 93 ++++--- .../camel/component/undertow/UndertowEndpoint.java | 40 +++ .../camel/component/undertow/UndertowHost.java | 19 ++ .../camel/component/undertow/UndertowProducer.java | 6 +- .../undertow/handlers/CamelWebSocketHandler.java | 285 +++++++++++++++++++-- .../CamelWebSocketHandlerSecuritySettingsTest.java | 63 +++++ .../ProviderWithServletWebSocketEndpointTest.java | 91 +++++++ .../spi/ProviderWithServletWebSocketTest.java | 66 +++++ ...ityProviderRolesFromComponentWebSocketTest.java | 78 ++++++ .../spi/SecurityProviderWebSocketTest.java | 160 ++++++++++++ .../spi/SecurityProviderWebSocketWrapTest.java | 89 +++++++ .../UndertowWsOAuthProfileConsumerWindowTest.java | 93 +++++-- .../ws/UndertowWsSecurityWithoutProviderTest.java | 197 ++++++++++++++ .../ROOT/pages/camel-4x-upgrade-guide-4_22.adoc | 18 ++ 19 files changed, 1408 insertions(+), 145 deletions(-) diff --git a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/undertow-component.adoc b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/undertow-component.adoc index 9ad58faee0ba..c35012a58843 100644 --- a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/undertow-component.adoc +++ b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/undertow-component.adoc @@ -125,6 +125,8 @@ the route is invoked. For WebSocket consumers, the token is validated during the HTTP upgrade handshake. The validation result is available on the `ONOPEN` event exchange when `fireWebSocketChannelEvents=true`, and on subsequent message exchanges for that connection; individual WebSocket messages are not revalidated. +While the consumer is stopped, upgrade requests to its path are still validated, and a connection whose +handshake was not validated does not receive the messages sent by producers on that path. ==== [NOTE] @@ -331,10 +333,26 @@ If there is an object passed to the component as parameter `securityConfiguratio Provider will be used for authentication of all requests. -Property `requireServletContext` of security providers forces the Undertow server to start -with servlet context. There will be no servlet actually handled. This feature is meant only +Property `requireServletContext` of security providers forces the Undertow server to run +with servlet context, as soon as an endpoint that uses such a provider, configured on the endpoint +or on the component, is started. There will be no servlet actually handled. This feature is meant only for use with servlet filters, which needs servlet context for their functionality. +==== WebSocket endpoints + +On WebSocket endpoints (`ws://` and `wss://`), the security provider and the `allowedRoles` option apply to the +upgrade request that opens a connection. A request that the provider rejects, or a request to an endpoint that has +`allowedRoles` but no security provider, is answered with the same status code as on an HTTP endpoint, and no +connection is opened. The headers that the provider adds are set on every exchange of the connection. + +All the endpoints of a WebSocket path share its connections, including the producers that send to its peers. The +security settings of the consumer apply to the whole path, and keep applying while the consumer is stopped. On a path +that only has producers, the settings of the producers apply. A connection that was opened before the current +settings applied, for example while the path only had producers without security settings, does not receive the +messages of the producers, and it is closed when it sends a message. + +The `handlers` and `accessLog` options apply to WebSocket consumers as well, before the upgrade. + == Examples === HTTP Producer Example diff --git a/components/camel-spring-parent/camel-undertow-spring-security/src/test/java/org/apache/camel/component/spring/security/SpringSecurityWebSocketTest.java b/components/camel-spring-parent/camel-undertow-spring-security/src/test/java/org/apache/camel/component/spring/security/SpringSecurityWebSocketTest.java new file mode 100644 index 000000000000..d05dedb8d0b0 --- /dev/null +++ b/components/camel-spring-parent/camel-undertow-spring-security/src/test/java/org/apache/camel/component/spring/security/SpringSecurityWebSocketTest.java @@ -0,0 +1,90 @@ +/* + * 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.camel.component.spring.security; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.WebSocket; +import java.net.http.WebSocketHandshakeException; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.TimeUnit; + +import io.undertow.util.StatusCodes; +import org.apache.camel.builder.RouteBuilder; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; + +/** + * The Spring Security provider applies to the upgrade request of WebSocket endpoints. + */ +class SpringSecurityWebSocketTest extends AbstractSpringSecurityBearerTokenTest { + + @Test + void allowedRoleConnects() throws Exception { + getMockFilter().setJwt(createToken("Alice", "user")); + + CompletableFuture<String> reply = new CompletableFuture<>(); + WebSocket webSocket = connect(new WebSocket.Listener() { + @Override + public void onOpen(WebSocket webSocket) { + webSocket.request(1); + } + + @Override + public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last) { + reply.complete(data.toString()); + return null; + } + }); + webSocket.sendText("hi", true).join(); + + assertEquals("Hello Alice!", reply.get(10, TimeUnit.SECONDS)); + webSocket.sendClose(WebSocket.NORMAL_CLOSURE, "done").join(); + } + + @Test + void otherRoleIsRefused() { + getMockFilter().setJwt(createToken("Tom", "wrongUser")); + + CompletionException thrown = assertThrows(CompletionException.class, () -> connect(new WebSocket.Listener() { + })); + WebSocketHandshakeException handshake = assertInstanceOf(WebSocketHandshakeException.class, thrown.getCause()); + assertEquals(StatusCodes.FORBIDDEN, handshake.getResponse().statusCode()); + } + + private WebSocket connect(WebSocket.Listener listener) { + return HttpClient.newHttpClient().newWebSocketBuilder() + .buildAsync(URI.create("ws://localhost:" + getPort() + "/myws"), listener) + .orTimeout(5, TimeUnit.SECONDS).join(); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + public void configure() { + from("undertow:ws://localhost:{{port}}/myws?allowedRoles=user") + .transform(simple("Hello ${in.header." + SpringSecurityProvider.PRINCIPAL_NAME_HEADER + "}!")) + .to("undertow:ws://localhost:{{port}}/myws"); + } + }; + } +} diff --git a/components/camel-undertow/src/main/docs/undertow-component.adoc b/components/camel-undertow/src/main/docs/undertow-component.adoc index 9ad58faee0ba..c35012a58843 100644 --- a/components/camel-undertow/src/main/docs/undertow-component.adoc +++ b/components/camel-undertow/src/main/docs/undertow-component.adoc @@ -125,6 +125,8 @@ the route is invoked. For WebSocket consumers, the token is validated during the HTTP upgrade handshake. The validation result is available on the `ONOPEN` event exchange when `fireWebSocketChannelEvents=true`, and on subsequent message exchanges for that connection; individual WebSocket messages are not revalidated. +While the consumer is stopped, upgrade requests to its path are still validated, and a connection whose +handshake was not validated does not receive the messages sent by producers on that path. ==== [NOTE] @@ -331,10 +333,26 @@ If there is an object passed to the component as parameter `securityConfiguratio Provider will be used for authentication of all requests. -Property `requireServletContext` of security providers forces the Undertow server to start -with servlet context. There will be no servlet actually handled. This feature is meant only +Property `requireServletContext` of security providers forces the Undertow server to run +with servlet context, as soon as an endpoint that uses such a provider, configured on the endpoint +or on the component, is started. There will be no servlet actually handled. This feature is meant only for use with servlet filters, which needs servlet context for their functionality. +==== WebSocket endpoints + +On WebSocket endpoints (`ws://` and `wss://`), the security provider and the `allowedRoles` option apply to the +upgrade request that opens a connection. A request that the provider rejects, or a request to an endpoint that has +`allowedRoles` but no security provider, is answered with the same status code as on an HTTP endpoint, and no +connection is opened. The headers that the provider adds are set on every exchange of the connection. + +All the endpoints of a WebSocket path share its connections, including the producers that send to its peers. The +security settings of the consumer apply to the whole path, and keep applying while the consumer is stopped. On a path +that only has producers, the settings of the producers apply. A connection that was opened before the current +settings applied, for example while the path only had producers without security settings, does not receive the +messages of the producers, and it is closed when it sends a message. + +The `handlers` and `accessLog` options apply to WebSocket consumers as well, before the upgrade. + == Examples === HTTP Producer Example diff --git a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/DefaultUndertowHost.java b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/DefaultUndertowHost.java index e8702dee5efd..1e33a6dd8b3f 100644 --- a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/DefaultUndertowHost.java +++ b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/DefaultUndertowHost.java @@ -49,6 +49,10 @@ public class DefaultUndertowHost implements UndertowHost { private Undertow undertow; private String hostString; private DeploymentManager deploymentManager; + // the rest or the root handler, depending on the first endpoint registered on the server + private HttpHandler serverHandler; + // the server handler, or the servlet deployment around it once an endpoint needs a servlet context + private volatile HttpHandler entryHandler; public DefaultUndertowHost(UndertowHostKey key) { this(key, null); @@ -70,6 +74,13 @@ public class DefaultUndertowHost implements UndertowHost { @Override public HttpHandler registerHandler( UndertowConsumer consumer, HttpHandlerRegistrationInfo registrationInfo, HttpHandler handler) { + return registerHandler(consumer != null ? consumer.getEndpoint() : null, consumer, registrationInfo, handler); + } + + @Override + public HttpHandler registerHandler( + UndertowEndpoint endpoint, UndertowConsumer consumer, HttpHandlerRegistrationInfo registrationInfo, + HttpHandler handler) { lock.lock(); try { if (undertow == null) { @@ -122,12 +133,10 @@ public class DefaultUndertowHost implements UndertowHost { } } - if (consumer != null && consumer.isRest()) { - // use the rest handler as its a rest consumer - undertow = registerHandler(consumer, builder, restHandler); - } else { - undertow = registerHandler(consumer, builder, rootHandler); - } + // use the rest handler as its a rest consumer + serverHandler = consumer != null && consumer.isRest() ? restHandler : rootHandler; + entryHandler = serverHandler; + undertow = builder.setHandler(exchange -> entryHandler.handleRequest(exchange)).build(); LOG.info("Starting Undertow server on {}://{}:{}", key.getSslContext() != null ? "https" : "http", key.getHost(), key.getPort()); @@ -151,6 +160,9 @@ public class DefaultUndertowHost implements UndertowHost { throw e; } } + if (deploymentManager == null && requiresServletContext(endpoint)) { + deployServletContext(); + } if (consumer != null && consumer.isRest()) { restHandler.addConsumer(consumer, handler); return restHandler; @@ -163,36 +175,44 @@ public class DefaultUndertowHost implements UndertowHost { } } - private Undertow registerHandler(UndertowConsumer consumer, Undertow.Builder builder, HttpHandler handler) { - UndertowSecurityProvider securityProvider = consumer == null - ? null - : consumer.getEndpoint().getComponent().getSecurityProvider() != null - ? consumer.getEndpoint().getComponent().getSecurityProvider() - : consumer.getEndpoint().getSecurityProvider(); - //if security provider needs servlet context, start empty servlet - if (securityProvider != null && securityProvider.requireServletContext()) { - DeploymentInfo deployment = Servlets.deployment() - .setContextPath("") - .setDisplayName("application") - .setDeploymentName("camel-undertow") - .setClassLoader(getClass().getClassLoader()) - //httpHandler for servlet is ignored, camel handler is used instead of it - .addOuterHandlerChainWrapper(h -> handler); - - deploymentManager = Servlets.newContainer().addDeployment(deployment); - deploymentManager.deploy(); - try { - return builder.setHandler(deploymentManager.start()).build(); - } catch (ServletException e) { - LOG.warn("Failed to start Undertow server on {}://{}:{}, reason: {}", - key.getSslContext() != null ? "https" : "http", key.getHost(), key.getPort(), e.getMessage()); - - throw new RuntimeException(e); - } - + /** + * Whether the security provider of the endpoint, or of its component, needs a servlet context, for example to run + * servlet filters. A producer endpoint can register first on a server, such as a WebSocket producer. + */ + private static boolean requiresServletContext(UndertowEndpoint endpoint) { + if (endpoint == null) { + return false; } + UndertowSecurityProvider endpointProvider = endpoint.getSecurityProvider(); + UndertowSecurityProvider componentProvider = endpoint.getComponent().getSecurityProvider(); + return endpointProvider != null && endpointProvider.requireServletContext() + || componentProvider != null && componentProvider.requireServletContext(); + } - return builder.setHandler(handler).build(); + /** + * Starts an empty servlet deployment around the handler of the server, so that every request has a servlet context. + */ + private void deployServletContext() { + HttpHandler handler = serverHandler; + DeploymentInfo deployment = Servlets.deployment() + .setContextPath("") + .setDisplayName("application") + .setDeploymentName("camel-undertow") + .setClassLoader(getClass().getClassLoader()) + //httpHandler for servlet is ignored, camel handler is used instead of it + .addOuterHandlerChainWrapper(h -> handler); + + DeploymentManager manager = Servlets.newContainer().addDeployment(deployment); + manager.deploy(); + try { + entryHandler = manager.start(); + } catch (ServletException e) { + LOG.warn("Failed to start the servlet context of the Undertow server on {}://{}:{}, reason: {}", + key.getSslContext() != null ? "https" : "http", key.getHost(), key.getPort(), e.getMessage()); + manager.undeploy(); + throw new RuntimeException(e); + } + deploymentManager = manager; } @Override @@ -212,11 +232,12 @@ public class DefaultUndertowHost implements UndertowHost { registrationInfo.isMatchOnUriPrefix()); stop = rootHandler.isEmpty(); } - if (deploymentManager != null) { - deploymentManager.undeploy(); - } - if (stop) { + // the servlet deployment serves every endpoint of the server, so it can only go with the server + if (deploymentManager != null) { + deploymentManager.undeploy(); + deploymentManager = null; + } LOG.info("Stopping Undertow server on {}://{}:{}", key.getSslContext() != null ? "https" : "http", key.getHost(), key.getPort()); diff --git a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowComponent.java b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowComponent.java index 67bf2ca62113..96c4d325f705 100644 --- a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowComponent.java +++ b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowComponent.java @@ -37,6 +37,7 @@ import org.apache.camel.Endpoint; import org.apache.camel.Processor; import org.apache.camel.Producer; import org.apache.camel.SSLContextParametersAware; +import org.apache.camel.component.undertow.handlers.CamelWebSocketHandler; import org.apache.camel.component.undertow.spi.UndertowSecurityProvider; import org.apache.camel.spi.Metadata; import org.apache.camel.spi.RestApiConsumerFactory; @@ -364,6 +365,18 @@ public class UndertowComponent extends DefaultComponent public HttpHandler registerEndpoint( UndertowConsumer consumer, HttpHandlerRegistrationInfo registrationInfo, SSLContext sslContext, HttpHandler handler) throws Exception { + return registerEndpoint(consumer != null ? consumer.getEndpoint() : null, consumer, registrationInfo, sslContext, + handler); + } + + /** + * Registers a handler on behalf of the given endpoint: the endpoint of the consumer, or a producer endpoint that + * registers a handler, such as a WebSocket producer, when {@code consumer} is {@code null}. + */ + public HttpHandler registerEndpoint( + UndertowEndpoint endpoint, UndertowConsumer consumer, HttpHandlerRegistrationInfo registrationInfo, + SSLContext sslContext, HttpHandler handler) + throws Exception { final URI uri = registrationInfo.getUri(); final UndertowHostKey key = new UndertowHostKey(uri.getHost(), uri.getPort(), sslContext); final UndertowHost host = undertowRegistry.computeIfAbsent(key, this::createUndertowHost); @@ -372,11 +385,18 @@ public class UndertowComponent extends DefaultComponent handlers.add(registrationInfo); HttpHandler handlerWrapped = handler; - if (this.securityProvider != null) { + if (handler instanceof CamelWebSocketHandler webSocketHandler) { + // the WebSocket handler of a path is shared by its consumer and producers, so it must stay registered as is. + // It is wrapped before the registration, so that it cannot receive a request unwrapped; when the path + // already has a handler, the registration keeps that one and this instance is not used + if (this.securityProvider != null) { + webSocketHandler.wrapWith(this.securityProvider); + } + } else if (this.securityProvider != null) { handlerWrapped = this.securityProvider.wrapHttpHandler(handler); } - return host.registerHandler(consumer, registrationInfo, handlerWrapped); + return host.registerHandler(endpoint, consumer, registrationInfo, handlerWrapped); } public void unregisterEndpoint( diff --git a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowConsumer.java b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowConsumer.java index 0cd70ca89b78..02c82e512bdf 100644 --- a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowConsumer.java +++ b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowConsumer.java @@ -21,7 +21,6 @@ import java.io.InputStream; import java.io.OutputStream; import java.net.URI; import java.nio.ByteBuffer; -import java.util.Arrays; import java.util.Collection; import java.util.List; import java.util.StringJoiner; @@ -99,11 +98,7 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su } public List<String> computeAllowedRoles() { - String allowedRolesString = getEndpoint().getAllowedRoles(); - if (allowedRolesString == null) { - allowedRolesString = getEndpoint().getComponent().getAllowedRoles(); - } - return allowedRolesString == null ? null : Arrays.asList(allowedRolesString.split("\\s*,\\s*")); + return getEndpoint().computeAllowedRoles(); } @Override @@ -118,7 +113,9 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su */ this.webSocketHandler = (CamelWebSocketHandler) endpoint.getComponent().registerEndpoint(this, endpoint.getHttpHandlerRegistrationInfo(), endpoint.getSslContext(), new CamelWebSocketHandler()); - this.webSocketHandler.setConsumer(this); + // the access log and the custom handlers run before the upgrade, as they run before an HTTP request + this.webSocketHandler.setConsumer(this, + wrapWithAccessLogAndHandlers(this.webSocketHandler.getUpgradeHandler(), endpoint)); } else { // allow for HTTP 1.1 continue HttpHandler httpHandler = new EagerFormParsingHandler().setNext(UndertowConsumer.this); @@ -127,22 +124,7 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su endpoint.getCamelContext(), endpoint.getOauthHttpSecurity(), httpHandler, endpoint.isOptionsEnabled(), this::isSuspended); } - if (endpoint.getAccessLog()) { - AccessLogReceiver accessLogReceiver; - if (endpoint.getAccessLogReceiver() != null) { - accessLogReceiver = endpoint.getAccessLogReceiver(); - } else { - accessLogReceiver = new JBossLoggingAccessLogReceiver(); - } - httpHandler = new AccessLogHandler( - httpHandler, - accessLogReceiver, - "common", - AccessLogHandler.class.getClassLoader()); - } - if (endpoint.getHandlers() != null) { - httpHandler = this.wrapHandler(httpHandler, endpoint); - } + httpHandler = wrapWithAccessLogAndHandlers(httpHandler, endpoint); endpoint.getComponent().registerEndpoint(this, endpoint.getHttpHandlerRegistrationInfo(), endpoint.getSslContext(), Handlers.httpContinueRead( // wrap with EagerFormParsingHandler to enable undertow form parsers @@ -150,6 +132,27 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su } } + private HttpHandler wrapWithAccessLogAndHandlers(HttpHandler handler, UndertowEndpoint endpoint) { + HttpHandler httpHandler = handler; + if (endpoint.getAccessLog()) { + AccessLogReceiver accessLogReceiver; + if (endpoint.getAccessLogReceiver() != null) { + accessLogReceiver = endpoint.getAccessLogReceiver(); + } else { + accessLogReceiver = new JBossLoggingAccessLogReceiver(); + } + httpHandler = new AccessLogHandler( + httpHandler, + accessLogReceiver, + "common", + AccessLogHandler.class.getClassLoader()); + } + if (endpoint.getHandlers() != null) { + httpHandler = this.wrapHandler(httpHandler, endpoint); + } + return httpHandler; + } + @Override protected void doStop() throws Exception { this.suspended = false; @@ -205,19 +208,9 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su return; } - if (getEndpoint().getSecurityProvider() != null) { - //security provider decides, whether endpoint is accessible - int statusCode = getEndpoint().getSecurityProvider().authenticate(httpExchange, computeAllowedRoles()); - if (statusCode != StatusCodes.OK) { - httpExchange.setStatusCode(statusCode); - httpExchange.endExchange(); - return; - } - } else if (computeAllowedRoles() != null && !computeAllowedRoles().isEmpty()) { - //this case could happen due to bad configuration - //if allowedRoles are present but securityProvider is not, access has to be denied in this case - LOG.warn("Illegal state caused by missing securitProvider but existing allowed roles!"); - httpExchange.setStatusCode(StatusCodes.FORBIDDEN); + int statusCode = getEndpoint().authenticate(httpExchange); + if (statusCode != StatusCodes.OK) { + httpExchange.setStatusCode(statusCode); httpExchange.endExchange(); return; } @@ -311,6 +304,7 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su final Exchange exchange = createExchange(true); // set header and body + setSecurityProviderHeaders(exchange.getIn(), channel); exchange.getIn().setHeader(UndertowConstants.CONNECTION_KEY, connectionKey); if (channel != null) { exchange.getIn().setHeader(UndertowConstants.CHANNEL, channel); @@ -339,6 +333,7 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su final Exchange exchange = createExchange(true); final Message in = exchange.getIn(); + setSecurityProviderHeaders(in, channel); in.setHeader(UndertowConstants.CONNECTION_KEY, connectionKey); in.setHeader(UndertowConstants.EVENT_TYPE, eventType.getCode()); in.setHeader(UndertowConstants.EVENT_TYPE_ENUM, eventType); @@ -412,28 +407,30 @@ public class UndertowConsumer extends DefaultConsumer implements HttpHandler, Su } /** - * Fails closed for WebSocket channels whose handshake was not OAuth validated while this consumer requires it. Such - * channels can exist when the shared {@link CamelWebSocketHandler} accepted a handshake while no consumer was set - * on it, for example between producer registration and consumer start, or across a consumer restart. Returns - * {@code true} when the event must not be delivered to the route; the channel is closed if still open. + * Fails closed for WebSocket channels whose handshake did not pass the security checks of this consumer's endpoint + * (OAuth validation, security provider, allowed roles and custom handlers). Such channels can exist when the shared + * {@link CamelWebSocketHandler} accepted a handshake before these checks applied to its path, for example while the + * path was only used by producers. Returns {@code true} when the event must not be delivered to the route; the + * channel is closed if still open. */ private boolean rejectUnauthenticatedWebSocketChannel(String connectionKey, WebSocketChannel channel) { - if (getEndpoint().getOauthHttpSecurity() == null) { - return false; - } - boolean authenticated = channel != null - && channel.getAttribute( - OAuthHttpSecuritySupport.OAUTH_TOKEN_VALIDATION_RESULT) instanceof OAuthTokenValidationResult; - if (authenticated) { + if (CamelWebSocketHandler.isAuthenticated(channel, getEndpoint())) { return false; } - LOG.warn("Rejecting WebSocket event from connection {} whose handshake was not OAuth validated", connectionKey); + LOG.warn("Rejecting WebSocket event from connection {} whose handshake did not pass the security checks", + connectionKey); if (channel != null && channel.isOpen()) { WebSockets.sendClose(CloseMessage.MSG_VIOLATES_POLICY, "Authentication required", channel, null); } return true; } + private static void setSecurityProviderHeaders(Message in, WebSocketChannel channel) { + if (channel != null) { + CamelWebSocketHandler.getSecurityProviderHeaders(channel).forEach(in::setHeader); + } + } + private static void setOAuthTokenValidationResult(Exchange exchange, WebSocketChannel channel) { if (channel == null) { return; diff --git a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowEndpoint.java b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowEndpoint.java index d97ae90edb26..24c23eecab29 100644 --- a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowEndpoint.java +++ b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowEndpoint.java @@ -17,6 +17,7 @@ package org.apache.camel.component.undertow; import java.net.URI; +import java.util.Arrays; import java.util.Iterator; import java.util.LinkedList; import java.util.List; @@ -26,7 +27,9 @@ import java.util.ServiceLoader; import javax.net.ssl.SSLContext; +import io.undertow.server.HttpServerExchange; import io.undertow.server.handlers.accesslog.AccessLogReceiver; +import io.undertow.util.StatusCodes; import org.apache.camel.AsyncEndpoint; import org.apache.camel.Category; import org.apache.camel.Consumer; @@ -487,6 +490,43 @@ public class UndertowEndpoint extends DefaultEndpoint this.allowedRoles = allowedRoles; } + /** + * The allowed roles of this endpoint, or of the component when the endpoint does not configure any. + */ + public List<String> computeAllowedRoles() { + String allowedRolesString = allowedRoles != null ? allowedRoles : getComponent().getAllowedRoles(); + return allowedRolesString == null ? null : Arrays.asList(allowedRolesString.split("\\s*,\\s*")); + } + + /** + * Whether requests to this endpoint are checked by the {@link UndertowSecurityProvider} or restricted to the + * allowed roles. + */ + public boolean requiresAuthentication() { + List<String> roles = computeAllowedRoles(); + return securityProvider != null || roles != null && !roles.isEmpty(); + } + + /** + * Applies the {@link UndertowSecurityProvider} and the allowed roles of this endpoint to a request. + * + * @return {@link StatusCodes#OK} if the request is allowed, otherwise the status code to reject it with + */ + public int authenticate(HttpServerExchange httpExchange) throws Exception { + List<String> roles = computeAllowedRoles(); + if (securityProvider != null) { + // security provider decides, whether endpoint is accessible + return securityProvider.authenticate(httpExchange, roles); + } + if (roles != null && !roles.isEmpty()) { + // this case could happen due to bad configuration + // if allowedRoles are present but securityProvider is not, access has to be denied in this case + LOG.warn("Illegal state caused by missing securityProvider but existing allowed roles!"); + return StatusCodes.FORBIDDEN; + } + return StatusCodes.OK; + } + @Override protected void doInit() throws Exception { super.doInit(); diff --git a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowHost.java b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowHost.java index 7559d1750b68..c9fa22883e24 100644 --- a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowHost.java +++ b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowHost.java @@ -46,6 +46,25 @@ public interface UndertowHost { */ HttpHandler registerHandler(UndertowConsumer consumer, HttpHandlerRegistrationInfo registrationInfo, HttpHandler handler); + /** + * Register a handler on behalf of the given endpoint, as + * {@link #registerHandler(UndertowConsumer, HttpHandlerRegistrationInfo, HttpHandler)} does. The endpoint is the + * endpoint of the consumer, or a producer endpoint that registers a handler, such as a WebSocket producer, when + * {@code consumer} is {@code null}. + * + * @param endpoint the endpoint that registers the handler + * @param consumer the consumer that registers the handler, or {@code null} for a producer + * @param registrationInfo the {@link HttpHandlerRegistrationInfo} related to {@code handler} + * @param handler the {@link HttpHandler} to register + * @return the given {@code handler} or a different {@link HttpHandler} that has been registered + * with the given {@link HttpHandlerRegistrationInfo} earlier. + */ + default HttpHandler registerHandler( + UndertowEndpoint endpoint, UndertowConsumer consumer, HttpHandlerRegistrationInfo registrationInfo, + HttpHandler handler) { + return registerHandler(consumer, registrationInfo, handler); + } + /** * Unregister a handler with the given {@link HttpHandlerRegistrationInfo}. Note that if * {@link #registerHandler(UndertowConsumer, HttpHandlerRegistrationInfo, HttpHandler)} was successfully invoked diff --git a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowProducer.java b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowProducer.java index 0d053cffb241..06a90683f806 100644 --- a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowProducer.java +++ b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/UndertowProducer.java @@ -249,8 +249,9 @@ public class UndertowProducer extends DefaultAsyncProducer { client = UndertowClient.getInstance(); if (endpoint.isWebSocket()) { - this.webSocketHandler = (CamelWebSocketHandler) endpoint.getComponent().registerEndpoint(null, + this.webSocketHandler = (CamelWebSocketHandler) endpoint.getComponent().registerEndpoint(endpoint, null, endpoint.getHttpHandlerRegistrationInfo(), endpoint.getSslContext(), new CamelWebSocketHandler()); + this.webSocketHandler.addProducer(endpoint); } LOG.debug("Created worker: {} with options: {}", worker, options); @@ -261,6 +262,9 @@ public class UndertowProducer extends DefaultAsyncProducer { super.doStop(); if (endpoint.isWebSocket()) { + if (webSocketHandler != null) { + webSocketHandler.removeProducer(endpoint); + } endpoint.getComponent().unregisterEndpoint(null, endpoint.getHttpHandlerRegistrationInfo(), endpoint.getSslContext()); } diff --git a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandler.java b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandler.java index 5b18d6d57e83..cf46381bde97 100644 --- a/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandler.java +++ b/components/camel-undertow/src/main/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandler.java @@ -24,13 +24,17 @@ import java.io.StringReader; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Collection; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Map.Entry; +import java.util.Objects; import java.util.Set; import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import java.util.function.Predicate; @@ -39,10 +43,12 @@ import java.util.stream.Collectors; import io.undertow.Handlers; import io.undertow.server.HttpHandler; import io.undertow.server.HttpServerExchange; +import io.undertow.server.handlers.ResponseCodeHandler; import io.undertow.util.AttachmentKey; import io.undertow.util.HeaderValues; import io.undertow.util.Headers; import io.undertow.util.MimeMappings; +import io.undertow.util.StatusCodes; import io.undertow.websockets.WebSocketConnectionCallback; import io.undertow.websockets.WebSocketProtocolHandshakeHandler; import io.undertow.websockets.core.AbstractReceiveListener; @@ -59,7 +65,9 @@ import org.apache.camel.RuntimeCamelException; import org.apache.camel.component.undertow.UndertowConstants; import org.apache.camel.component.undertow.UndertowConstants.EventType; import org.apache.camel.component.undertow.UndertowConsumer; +import org.apache.camel.component.undertow.UndertowEndpoint; import org.apache.camel.component.undertow.UndertowProducer; +import org.apache.camel.component.undertow.spi.UndertowSecurityProvider; import org.apache.camel.converter.IOConverter; import org.apache.camel.http.base.OAuthHttpSecuritySupport; import org.apache.camel.http.base.OAuthHttpSecuritySupport.Validation; @@ -78,11 +86,24 @@ public class CamelWebSocketHandler implements HttpHandler { // WebSocket handshakes use WebSocketHttpExchange attachments; HTTP requests use UndertowConsumer's key. private static final AttachmentKey<OAuthTokenValidationResult> OAUTH_TOKEN_VALIDATION_RESULT_ATTACHMENT = AttachmentKey.create(OAuthTokenValidationResult.class); + private static final AttachmentKey<HandshakeResult> HANDSHAKE_RESULT_ATTACHMENT + = AttachmentKey.create(HandshakeResult.class); + private static final String HANDSHAKE_RESULT = CamelWebSocketHandler.class.getName() + ".handshakeResult"; + private static final AttachmentKey<String> CONSUMER_HANDLERS_ATTACHMENT = AttachmentKey.create(String.class); private final UndertowWebSocketConnectionCallback callback; private UndertowConsumer consumer; + /** + * The endpoint of the last consumer set on this handler, whose security settings apply to the path. + */ + private UndertowEndpoint consumerEndpoint; + + private final List<UndertowEndpoint> producerEndpoints = new CopyOnWriteArrayList<>(); + + private final Set<String> warnedProducerEndpoints = ConcurrentHashMap.newKeySet(); + private final Lock consumerLock = new ReentrantLock(); private final WebSocketProtocolHandshakeHandler delegate; @@ -91,6 +112,17 @@ public class CamelWebSocketHandler implements HttpHandler { private final UndertowReceiveListener receiveListener; + private final HttpHandler upgradeHandler = this::upgrade; + + private final HttpHandler consumerRequestHandler = this::handleConsumerRequest; + + /** + * The handlers of the consumer, such as its access log, followed by {@link #upgradeHandler}. + */ + private volatile ConsumerHandlers consumerHandlers = new ConsumerHandlers(null, upgradeHandler); + + private volatile HttpHandler entryHandler = consumerRequestHandler; + public CamelWebSocketHandler() { this.receiveListener = new UndertowReceiveListener(); this.callback = new UndertowWebSocketConnectionCallback(); @@ -142,37 +174,144 @@ public class CamelWebSocketHandler implements HttpHandler { */ @Override public void handleRequest(HttpServerExchange exchange) throws Exception { - UndertowConsumer currentConsumer = getConsumer(); - if (currentConsumer != null && currentConsumer.getEndpoint().getOauthHttpSecurity() != null) { - if (exchange.isInIoThread()) { - exchange.dispatch(this); - return; + entryHandler.handleRequest(exchange); + } + + /** + * The handler that applies the security settings of the path to an upgrade request and then performs the WebSocket + * handshake. A consumer can run its own handlers, such as its access log, before it. + */ + public HttpHandler getUpgradeHandler() { + return upgradeHandler; + } + + /** + * Lets the given security provider wrap this handler, as it wraps the handlers of HTTP endpoints. The handler stays + * registered as is, so that it remains shared by the consumer and the producers of the path. + */ + public void wrapWith(UndertowSecurityProvider securityProvider) throws Exception { + HttpHandler wrapped = securityProvider.wrapHttpHandler(consumerRequestHandler); + // a provider that returns no handler disables the path, as it does for HTTP endpoints + this.entryHandler = wrapped != null ? wrapped : ResponseCodeHandler.HANDLE_405; + } + + private void handleConsumerRequest(HttpServerExchange exchange) throws Exception { + ConsumerHandlers handlers = consumerHandlers; + if (handlers.endpoint != null) { + // the request only reaches the upgrade if the handlers of this consumer let it through + exchange.putAttachment(CONSUMER_HANDLERS_ATTACHMENT, handlers.endpoint.getEndpointUri()); + } + handlers.handler.handleRequest(exchange); + } + + private void upgrade(HttpServerExchange exchange) throws Exception { + List<UndertowEndpoint> endpoints = securityEndpoints(); + if (endpoints.stream().noneMatch(CamelWebSocketHandler::hasSecurityChecks)) { + this.delegate.handleRequest(exchange); + return; + } + if (exchange.isInIoThread()) { + exchange.dispatch(upgradeHandler); + return; + } + Set<String> authenticatedEndpoints = new HashSet<>(); + Map<String, Object> headers = new HashMap<>(); + for (UndertowEndpoint endpoint : endpoints) { + OAuthHttpSecuritySupport oauthHttpSecurity = endpoint.getOauthHttpSecurity(); + if (oauthHttpSecurity != null) { + Validation validation = oauthHttpSecurity.validate(endpoint.getCamelContext(), authorizationHeaders(exchange)); + exchange.getRequestHeaders().remove(Headers.AUTHORIZATION); + if (!validation.isAuthenticated()) { + reject(exchange, validation); + return; + } + exchange.putAttachment(OAUTH_TOKEN_VALIDATION_RESULT_ATTACHMENT, validation.getValidationResult()); } - OAuthHttpSecuritySupport oauthHttpSecurity = currentConsumer.getEndpoint().getOauthHttpSecurity(); - Validation validation = oauthHttpSecurity.validate(currentConsumer.getEndpoint().getCamelContext(), - authorizationHeaders(exchange)); - exchange.getRequestHeaders().remove(Headers.AUTHORIZATION); - if (!validation.isAuthenticated()) { - reject(exchange, validation); - return; + if (endpoint.requiresAuthentication()) { + int statusCode = endpoint.authenticate(exchange); + if (statusCode != StatusCodes.OK) { + exchange.setStatusCode(statusCode); + exchange.endExchange(); + return; + } + if (endpoint.getSecurityProvider() != null) { + endpoint.getSecurityProvider().addHeader(headers::put, exchange); + } } - exchange.putAttachment(OAUTH_TOKEN_VALIDATION_RESULT_ATTACHMENT, validation.getValidationResult()); + if (requiresHandshakeResult(endpoint) && (endpoint.getHandlers() == null + || endpoint.getEndpointUri().equals(exchange.getAttachment(CONSUMER_HANDLERS_ATTACHMENT)))) { + authenticatedEndpoints.add(endpoint.getEndpointUri()); + } + } + if (!authenticatedEndpoints.isEmpty()) { + exchange.putAttachment(HANDSHAKE_RESULT_ATTACHMENT, + new HandshakeResult(authenticatedEndpoints, headers)); } this.delegate.handleRequest(exchange); } - private UndertowConsumer getConsumer() { + /** + * The endpoints whose security settings apply to the path: the consumer's, also while it is stopped, otherwise the + * producers'. All the Camel endpoints of a path share its WebSocket connections. + */ + private List<UndertowEndpoint> securityEndpoints() { consumerLock.lock(); try { - return consumer; + if (consumerEndpoint != null) { + return List.of(consumerEndpoint); + } } finally { consumerLock.unlock(); } + return producerEndpoints.stream().distinct().toList(); + } + + private static boolean hasSecurityChecks(UndertowEndpoint endpoint) { + return endpoint.getOauthHttpSecurity() != null || requiresHandshakeResult(endpoint); + } + + /** + * Whether a channel must have passed the security provider, allowed roles or custom handlers of the endpoint during + * its handshake. + */ + private static boolean requiresHandshakeResult(UndertowEndpoint endpoint) { + return endpoint.requiresAuthentication() || endpoint.getHandlers() != null; + } + + private static boolean isAuthenticated(WebSocketChannel channel, List<UndertowEndpoint> endpoints) { + for (UndertowEndpoint endpoint : endpoints) { + if (!isAuthenticated(channel, endpoint)) { + return false; + } + } + return true; + } + + /** + * Whether the handshake of the given channel passed the security checks of the given endpoint: its OAuth + * validation, its security provider, its allowed roles and its custom handlers. Always {@code true} for an endpoint + * without such checks. + */ + public static boolean isAuthenticated(WebSocketChannel channel, UndertowEndpoint endpoint) { + if (endpoint.getOauthHttpSecurity() != null && (channel == null + || !(channel.getAttribute( + OAuthHttpSecuritySupport.OAUTH_TOKEN_VALIDATION_RESULT) instanceof OAuthTokenValidationResult))) { + return false; + } + if (requiresHandshakeResult(endpoint)) { + return channel != null + && channel.getAttribute(HANDSHAKE_RESULT) instanceof HandshakeResult result + && result.endpointUris.contains(endpoint.getEndpointUri()); + } + return true; } - private boolean requiresOAuth() { - UndertowConsumer currentConsumer = getConsumer(); - return currentConsumer != null && currentConsumer.getEndpoint().getOauthHttpSecurity() != null; + /** + * The headers that the security providers added during the handshake of the given channel. + */ + public static Map<String, Object> getSecurityProviderHeaders(WebSocketChannel channel) { + return channel.getAttribute(HANDSHAKE_RESULT) instanceof HandshakeResult result + ? result.headers : Collections.emptyMap(); } private static void reject(HttpServerExchange exchange, Validation validation) { @@ -213,8 +352,12 @@ public class CamelWebSocketHandler implements HttpHandler { Predicate<WebSocketChannel> peerFilter, Object message, final int timeout, final Exchange camelExchange, final AsyncCallback camelCallback) throws IOException { - List<WebSocketChannel> targetPeers - = delegate.getPeerConnections().stream().filter(peerFilter).collect(Collectors.toList()); + // only the peers whose handshake passed the security checks of the path receive messages + List<UndertowEndpoint> endpoints = securityEndpoints(); + List<WebSocketChannel> targetPeers = delegate.getPeerConnections().stream() + .filter(peer -> isAuthenticated(peer, endpoints)) + .filter(peerFilter) + .collect(Collectors.toList()); if (targetPeers.isEmpty()) { camelCallback.done(true); return true; @@ -232,6 +375,15 @@ public class CamelWebSocketHandler implements HttpHandler { * @param consumer the {@link UndertowConsumer} to set */ public void setConsumer(UndertowConsumer consumer) { + setConsumer(consumer, null); + } + + /** + * @param consumer the {@link UndertowConsumer} to set + * @param consumerHandler the handler that runs before {@link #getUpgradeHandler()} for this consumer, or + * {@code null} + */ + public void setConsumer(UndertowConsumer consumer, HttpHandler consumerHandler) { consumerLock.lock(); try { if (consumer != null && this.consumer != null) { @@ -240,11 +392,59 @@ public class CamelWebSocketHandler implements HttpHandler { + ".setConsumer(UndertowConsumer) with a non-null consumer before unsetting it via setConsumer(null)"); } this.consumer = consumer; + if (consumer != null) { + // both are kept when the consumer is unset, so that the path stays guarded while the consumer is stopped + this.consumerEndpoint = consumer.getEndpoint(); + this.consumerHandlers + = new ConsumerHandlers( + consumer.getEndpoint(), consumerHandler != null ? consumerHandler : upgradeHandler); + producerEndpoints.forEach(this::warnIfProducerSettingsUnused); + } } finally { consumerLock.unlock(); } } + /** + * Registers the endpoint of a producer that sends to the peers of this handler. A path without a consumer applies + * the security settings of its producers. + */ + public void addProducer(UndertowEndpoint endpoint) { + producerEndpoints.add(endpoint); + warnIfProducerSettingsUnused(endpoint); + } + + private void warnIfProducerSettingsUnused(UndertowEndpoint producerEndpoint) { + UndertowEndpoint pathConsumerEndpoint; + consumerLock.lock(); + try { + pathConsumerEndpoint = consumerEndpoint; + } finally { + consumerLock.unlock(); + } + if (pathConsumerEndpoint != null && hasUnusedSecuritySettings(pathConsumerEndpoint, producerEndpoint) + && warnedProducerEndpoints.add(producerEndpoint.getEndpointUri())) { + LOG.warn("The security settings of {} are not used: the settings of the consumer {} apply to its WebSocket path", + producerEndpoint, pathConsumerEndpoint); + } + } + + /** + * Whether the producer endpoint has security settings of its own that differ from the ones of the consumer + * endpoint, which apply to the path instead. + */ + static boolean hasUnusedSecuritySettings(UndertowEndpoint consumerEndpoint, UndertowEndpoint producerEndpoint) { + boolean ownSettings = producerEndpoint.getAllowedRoles() != null + || producerEndpoint.getSecurityConfiguration() != null + || producerEndpoint.getSecurityProvider() != producerEndpoint.getComponent().getSecurityProvider(); + return ownSettings && (producerEndpoint.getSecurityProvider() != consumerEndpoint.getSecurityProvider() + || !Objects.equals(producerEndpoint.computeAllowedRoles(), consumerEndpoint.computeAllowedRoles())); + } + + public void removeProducer(UndertowEndpoint endpoint) { + producerEndpoints.remove(endpoint); + } + void sendEventNotificationIfNeeded( String connectionKey, WebSocketHttpExchange transportExchange, WebSocketChannel channel, EventType eventType) { consumerLock.lock(); @@ -427,6 +627,33 @@ public class CamelWebSocketHandler implements HttpHandler { } + /** + * The handlers of a consumer, which end with {@link #upgradeHandler}, and the endpoint of that consumer. + */ + private static final class ConsumerHandlers { + private final UndertowEndpoint endpoint; + private final HttpHandler handler; + + private ConsumerHandlers(UndertowEndpoint endpoint, HttpHandler handler) { + this.endpoint = endpoint; + this.handler = handler; + } + } + + /** + * The endpoints whose security provider, allowed roles and custom handlers a handshake passed, and the headers that + * the providers added. + */ + private static final class HandshakeResult { + private final Set<String> endpointUris; + private final Map<String, Object> headers; + + private HandshakeResult(Set<String> endpointUris, Map<String, Object> headers) { + this.endpointUris = Collections.unmodifiableSet(endpointUris); + this.headers = Collections.unmodifiableMap(headers); + } + } + /** * Sets the {@link UndertowReceiveListener} to the given channel on connect. */ @@ -440,18 +667,22 @@ public class CamelWebSocketHandler implements HttpHandler { LOG.trace("onConnect {}", exchange); OAuthTokenValidationResult oauthTokenValidationResult = exchange.getAttachment(OAUTH_TOKEN_VALIDATION_RESULT_ATTACHMENT); - if (oauthTokenValidationResult == null && requiresOAuth()) { - // the handshake completed without token validation, for example while no consumer was set on - // this handler, but the current consumer requires it: fail closed - LOG.warn("Closing WebSocket channel whose handshake was not OAuth validated"); + if (oauthTokenValidationResult != null) { + channel.setAttribute(OAuthHttpSecuritySupport.OAUTH_TOKEN_VALIDATION_RESULT, oauthTokenValidationResult); + } + HandshakeResult handshakeResult = exchange.getAttachment(HANDSHAKE_RESULT_ATTACHMENT); + if (handshakeResult != null) { + channel.setAttribute(HANDSHAKE_RESULT, handshakeResult); + } + if (!isAuthenticated(channel, securityEndpoints())) { + // the handshake did not pass the security checks that now apply to the path, for example because it + // completed before a consumer requiring them was set on this handler: fail closed + LOG.warn("Closing WebSocket channel whose handshake did not pass the security checks"); WebSockets.sendClose(CloseMessage.MSG_VIOLATES_POLICY, "Authentication required", channel, null); return; } final String connectionKey = UUID.randomUUID().toString(); channel.setAttribute(UndertowConstants.CONNECTION_KEY, connectionKey); - if (oauthTokenValidationResult != null) { - channel.setAttribute(OAuthHttpSecuritySupport.OAUTH_TOKEN_VALIDATION_RESULT, oauthTokenValidationResult); - } channel.getReceiveSetter().set(receiveListener); channel.addCloseTask(closeListener); sendEventNotificationIfNeeded(connectionKey, exchange, channel, EventType.ONOPEN); diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandlerSecuritySettingsTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandlerSecuritySettingsTest.java new file mode 100644 index 000000000000..6e26a728f74a --- /dev/null +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandlerSecuritySettingsTest.java @@ -0,0 +1,63 @@ +/* + * 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.camel.component.undertow.handlers; + +import org.apache.camel.CamelContext; +import org.apache.camel.component.undertow.UndertowEndpoint; +import org.apache.camel.impl.DefaultCamelContext; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * The security settings that a WebSocket producer configures for itself are not used when a consumer is on the path, + * which is reported when they differ from the settings of the consumer. + */ +class CamelWebSocketHandlerSecuritySettingsTest { + + @Test + void producerWithoutItsOwnSettings() throws Exception { + assertFalse(hasUnusedSecuritySettings("allowedRoles=user", "sendToAll=true")); + } + + @Test + void producerWithTheSameSettingsAsTheConsumer() throws Exception { + assertFalse(hasUnusedSecuritySettings("allowedRoles=admin", "allowedRoles=admin")); + } + + @Test + void producerWithOtherRolesThanTheConsumer() throws Exception { + assertTrue(hasUnusedSecuritySettings("allowedRoles=user", "allowedRoles=admin")); + } + + @Test + void producerWithRolesOnAPathWhoseConsumerHasNone() throws Exception { + assertTrue(hasUnusedSecuritySettings("fireWebSocketChannelEvents=true", "allowedRoles=admin")); + } + + private static boolean hasUnusedSecuritySettings(String consumerOptions, String producerOptions) throws Exception { + try (CamelContext context = new DefaultCamelContext()) { + context.start(); + UndertowEndpoint consumerEndpoint + = context.getEndpoint("undertow:ws://localhost:8080/path?" + consumerOptions, UndertowEndpoint.class); + UndertowEndpoint producerEndpoint + = context.getEndpoint("undertow:ws://localhost:8080/path?" + producerOptions, UndertowEndpoint.class); + return CamelWebSocketHandler.hasUnusedSecuritySettings(consumerEndpoint, producerEndpoint); + } + } +} diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/ProviderWithServletWebSocketEndpointTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/ProviderWithServletWebSocketEndpointTest.java new file mode 100644 index 000000000000..2cdffab50254 --- /dev/null +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/ProviderWithServletWebSocketEndpointTest.java @@ -0,0 +1,91 @@ +/* + * 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.camel.component.undertow.spi; + +import java.util.concurrent.TimeUnit; + +import org.apache.camel.CamelContext; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.undertow.BaseUndertowTest; +import org.apache.camel.test.infra.common.http.WebsocketTestClient; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * A security provider configured on a WebSocket endpoint that requires the servlet context gets it, whichever endpoint + * registers first on the port. + */ +class ProviderWithServletWebSocketEndpointTest extends BaseUndertowTest { + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext context = super.createCamelContext(); + context.getRegistry().bind("servletProvider", new ProviderWithServletTest.MockSecurityProvider()); + return context; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + // a path only used by a producer, which configures the provider + from("direct:feed") + .to("undertow:ws://localhost:{{port}}/feed?securityProvider=#servletProvider&sendToAll=true"); + + // the producer of the route registers first on the port, without a provider + from("undertow:ws://localhost:{{port2}}/chat?securityProvider=#servletProvider") + .to("mock:chat") + .transform(simple("${in.header." + AbstractSecurityProviderTest.PRINCIPAL_PARAMETER + "}")) + .to("undertow:ws://localhost:{{port2}}/chat"); + } + }; + } + + @Test + void producerOnlyPathWithItsOwnProvider() { + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/feed"); + client.connect(); + // the server registers the connection shortly after the client has completed the upgrade + await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + template.sendBody("direct:feed", "update"); + assertFalse(client.getReceived().isEmpty()); + }); + + assertEquals("update", client.getReceived(String.class).get(0)); + client.close(); + } + + @Test + void consumerWithItsOwnProviderOnAPortWhereAProducerRegisteredFirst() throws Exception { + getMockEndpoint("mock:chat").expectedBodiesReceived("hello"); + + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort2() + "/chat", 1); + client.connect(); + client.sendTextMessage("hello"); + + MockEndpoint.assertIsSatisfied(context); + assertTrue(client.await(10)); + assertEquals("user", client.getReceived(String.class).get(0)); + client.close(); + } +} diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/ProviderWithServletWebSocketTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/ProviderWithServletWebSocketTest.java new file mode 100644 index 000000000000..b04f8f783fed --- /dev/null +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/ProviderWithServletWebSocketTest.java @@ -0,0 +1,66 @@ +/* + * 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.camel.component.undertow.spi; + +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.test.infra.common.http.WebsocketTestClient; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * A security provider that requires the servlet context gets it for WebSocket upgrade requests too, also when the first + * endpoint registered on the port is a WebSocket producer. + */ +class ProviderWithServletWebSocketTest extends AbstractProviderServletTest { + + @BeforeAll + static void initProvider() throws Exception { + createSecurtyProviderConfigurationFile(ProviderWithServletTest.MockSecurityProvider.class); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + // the producer of the route is started, and registered on the port, before its consumer + from("undertow:ws://localhost:{{port}}/foo?allowedRoles=user") + .to("mock:input") + .transform(simple("${in.header." + AbstractSecurityProviderTest.PRINCIPAL_PARAMETER + "}")) + .to("undertow:ws://localhost:{{port}}/foo"); + } + }; + } + + @Test + void upgradeRequestHasTheServletContext() throws Exception { + getMockEndpoint("mock:input").expectedBodiesReceived("hello"); + + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/foo", 1); + client.connect(); + client.sendTextMessage("hello"); + + MockEndpoint.assertIsSatisfied(context); + assertTrue(client.await(10)); + assertEquals("user", client.getReceived(String.class).get(0)); + client.close(); + } +} diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderRolesFromComponentWebSocketTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderRolesFromComponentWebSocketTest.java new file mode 100644 index 000000000000..c46ec37b0852 --- /dev/null +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderRolesFromComponentWebSocketTest.java @@ -0,0 +1,78 @@ +/* + * 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.camel.component.undertow.spi; + +import java.net.http.WebSocketHandshakeException; +import java.util.concurrent.CompletionException; + +import io.undertow.util.StatusCodes; +import org.apache.camel.CamelContext; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.undertow.UndertowComponent; +import org.apache.camel.test.infra.common.http.WebsocketTestClient; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; + +/** + * The allowed roles of the component apply to the WebSocket endpoints that do not configure their own. + */ +class SecurityProviderRolesFromComponentWebSocketTest extends AbstractSecurityProviderTest { + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext camelContext = super.createCamelContext(); + camelContext.getComponent("undertow", UndertowComponent.class).setAllowedRoles("user"); + return camelContext; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("undertow:ws://localhost:{{port}}/roles").to("mock:input"); + } + }; + } + + @Test + void roleOfTheComponentIsAllowed() throws Exception { + securityConfiguration.setRoleToAssign("user"); + getMockEndpoint("mock:input").expectedBodiesReceived("hello"); + + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/roles"); + client.connect(); + client.sendTextMessage("hello"); + + MockEndpoint.assertIsSatisfied(context); + client.close(); + } + + @Test + void otherRoleIsRefused() { + securityConfiguration.setRoleToAssign("admin"); + + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/roles"); + CompletionException thrown = assertThrows(CompletionException.class, client::connect); + WebSocketHandshakeException handshake = assertInstanceOf(WebSocketHandshakeException.class, thrown.getCause()); + assertEquals(StatusCodes.FORBIDDEN, handshake.getResponse().statusCode()); + } +} diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderWebSocketTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderWebSocketTest.java new file mode 100644 index 000000000000..ad96dd162392 --- /dev/null +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderWebSocketTest.java @@ -0,0 +1,160 @@ +/* + * 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.camel.component.undertow.spi; + +import java.net.http.WebSocketHandshakeException; +import java.util.concurrent.CompletionException; +import java.util.concurrent.TimeUnit; + +import io.undertow.util.StatusCodes; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.test.infra.common.http.WebsocketTestClient; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * The security provider and the allowed roles apply to the upgrade request of WebSocket endpoints: on a path with a + * consumer, on a path only used by producers, and on a path whose consumer is stopped. + */ +class SecurityProviderWebSocketTest extends AbstractSecurityProviderTest { + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("undertow:ws://localhost:{{port}}/wssecure?allowedRoles=user") + .to("mock:input") + .transform(simple("echo:${body}")) + .to("undertow:ws://localhost:{{port}}/wssecure?sendToAll=true"); + + from("direct:feed") + .to("undertow:ws://localhost:{{port}}/feed?sendToAll=true&allowedRoles=user"); + + from("direct:stopped") + .to("undertow:ws://localhost:{{port}}/stopped?sendToAll=true"); + from("undertow:ws://localhost:{{port}}/stopped?allowedRoles=user").routeId("stopped") + .to("mock:stopped"); + + from("undertow:ws://localhost:{{port}}/prefix?matchOnUriPrefix=true&allowedRoles=user") + .to("mock:prefix"); + } + }; + } + + @Test + void subPathOfAPrefixPathIsGuarded() throws Exception { + securityConfiguration.setRoleToAssign("admin"); + assertRefused("/prefix/sub"); + + securityConfiguration.setRoleToAssign("user"); + MockEndpoint prefix = getMockEndpoint("mock:prefix"); + prefix.expectedBodiesReceived("hello"); + + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/prefix/sub"); + client.connect(); + client.sendTextMessage("hello"); + + prefix.assertIsSatisfied(); + client.close(); + } + + @Test + void matchingRoleConnectsAndExchangesMessages() throws Exception { + securityConfiguration.setRoleToAssign("user"); + MockEndpoint input = getMockEndpoint("mock:input"); + input.expectedBodiesReceived("ping"); + // the header added by the security provider during the upgrade + input.expectedHeaderReceived(PRINCIPAL_PARAMETER, "user"); + + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/wssecure", 1); + client.connect(); + client.sendTextMessage("ping"); + + input.assertIsSatisfied(); + assertTrue(client.await(10)); + assertEquals("echo:ping", client.getReceived(String.class).get(0)); + client.close(); + } + + @Test + void mismatchedRoleIsRefused() { + securityConfiguration.setRoleToAssign("admin"); + + assertRefused("/wssecure"); + } + + @Test + void missingRoleIsRefused() { + securityConfiguration.setRoleToAssign(null); + + assertRefused("/wssecure"); + } + + @Test + void producerOnlyPathAppliesTheProducerSettings() throws Exception { + securityConfiguration.setRoleToAssign("admin"); + assertRefused("/feed"); + + securityConfiguration.setRoleToAssign("user"); + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/feed"); + client.connect(); + // the server registers the connection shortly after the client has completed the upgrade + await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> { + template.sendBody("direct:feed", "update"); + assertFalse(client.getReceived().isEmpty()); + }); + + assertEquals("update", client.getReceived(String.class).get(0)); + client.close(); + } + + @Test + void pathStaysGuardedWhileTheConsumerIsStopped() throws Exception { + context.getRouteController().stopRoute("stopped"); + + securityConfiguration.setRoleToAssign("admin"); + assertRefused("/stopped"); + + securityConfiguration.setRoleToAssign("user"); + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + "/stopped"); + client.connect(); + context.getRouteController().startRoute("stopped"); + + MockEndpoint stopped = getMockEndpoint("mock:stopped"); + stopped.expectedBodiesReceived("hello"); + client.sendTextMessage("hello"); + + stopped.assertIsSatisfied(); + client.close(); + } + + private void assertRefused(String path) { + WebsocketTestClient client = new WebsocketTestClient("ws://localhost:" + getPort() + path); + + CompletionException thrown = assertThrows(CompletionException.class, client::connect); + WebSocketHandshakeException handshake = assertInstanceOf(WebSocketHandshakeException.class, thrown.getCause()); + assertEquals(StatusCodes.FORBIDDEN, handshake.getResponse().statusCode()); + } +} diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderWebSocketWrapTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderWebSocketWrapTest.java new file mode 100644 index 000000000000..7c7709a68b92 --- /dev/null +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/spi/SecurityProviderWebSocketWrapTest.java @@ -0,0 +1,89 @@ +/* + * 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.camel.component.undertow.spi; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.WebSocket; +import java.net.http.WebSocketHandshakeException; +import java.util.concurrent.CompletionException; +import java.util.concurrent.TimeUnit; + +import io.undertow.util.StatusCodes; +import org.apache.camel.CamelContext; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.junit.jupiter.api.Test; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; + +/** + * A security provider that wraps the HTTP handlers wraps the WebSocket endpoints as well. + */ +class SecurityProviderWebSocketWrapTest extends AbstractSecurityProviderTest { + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext camelContext = super.createCamelContext(); + securityConfiguration.setWrapHttpHandler(next -> exchange -> { + if ("yes".equals(exchange.getRequestHeaders().getFirst("X-Allow"))) { + next.handleRequest(exchange); + } else { + exchange.setStatusCode(StatusCodes.UNAUTHORIZED); + exchange.endExchange(); + } + }); + return camelContext; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("undertow:ws://localhost:{{port}}/wrapped?allowedRoles=user").to("mock:wrapped"); + } + }; + } + + @Test + void wrapperAppliesToTheUpgrade() throws Exception { + securityConfiguration.setRoleToAssign("user"); + + CompletionException thrown = assertThrows(CompletionException.class, () -> connect(null)); + WebSocketHandshakeException handshake = assertInstanceOf(WebSocketHandshakeException.class, thrown.getCause()); + assertEquals(StatusCodes.UNAUTHORIZED, handshake.getResponse().statusCode()); + + getMockEndpoint("mock:wrapped").expectedBodiesReceived("hello"); + WebSocket webSocket = connect("yes"); + webSocket.sendText("hello", true).join(); + + MockEndpoint.assertIsSatisfied(context); + webSocket.sendClose(WebSocket.NORMAL_CLOSURE, "done").join(); + } + + private WebSocket connect(String allow) { + WebSocket.Builder builder = HttpClient.newHttpClient().newWebSocketBuilder(); + if (allow != null) { + builder.header("X-Allow", allow); + } + return builder.buildAsync(URI.create("ws://localhost:" + getPort() + "/wrapped"), new WebSocket.Listener() { + }).orTimeout(5, TimeUnit.SECONDS).join(); + } +} diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsOAuthProfileConsumerWindowTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsOAuthProfileConsumerWindowTest.java index 9c00b6166c46..06a19be0e1bc 100644 --- a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsOAuthProfileConsumerWindowTest.java +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsOAuthProfileConsumerWindowTest.java @@ -19,8 +19,12 @@ package org.apache.camel.component.undertow.ws; import java.net.URI; import java.net.http.HttpClient; import java.net.http.WebSocket; +import java.net.http.WebSocketHandshakeException; +import java.util.List; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.CompletionStage; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @@ -33,11 +37,15 @@ import org.apache.camel.spi.OAuthTokenValidationFactory; import org.junit.jupiter.api.Test; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; /** - * Verifies that WebSocket channels whose handshake completed while no consumer was set on the shared - * {@code CamelWebSocketHandler} (here: while the secured consumer route was stopped and a producer route kept the - * handler registered) cannot deliver messages to the route once the secured consumer is back. + * Verifies the OAuth validation of a WebSocket path whose secured consumer is not running while a producer route keeps + * the shared {@code CamelWebSocketHandler} registered: while the consumer is stopped its settings keep applying to the + * path, and a channel whose handshake completed before the consumer was first started is neither served by the route + * nor sent messages once the secured consumer is up. */ public class UndertowWsOAuthProfileConsumerWindowTest extends BaseUndertowTest { @@ -51,37 +59,47 @@ public class UndertowWsOAuthProfileConsumerWindowTest extends BaseUndertowTest { } @Test - public void closesChannelEstablishedWhileConsumerWasStopped() throws Exception { - // stop the secured consumer; the producer route on the same path keeps the shared WebSocket handler registered + public void closesChannelEstablishedBeforeSecuredConsumerStarted() throws Exception { + // the secured consumer has not been started yet; the producer route on the same path keeps the shared + // WebSocket handler registered, so the handshake completes without OAuth validation + RecordingListener listener = new RecordingListener(); + WebSocket webSocket = connect(null, listener); + + // once the secured consumer is up, the unauthenticated channel must neither get messages nor reach the route + context.getRouteController().startRoute("secureWs"); + template.sendBody("direct:keepalive", "broadcast"); + + webSocket.sendText("hello", true).orTimeout(5, TimeUnit.SECONDS).join(); + + assertEquals(CloseMessage.MSG_VIOLATES_POLICY, listener.closeCode.orTimeout(10, TimeUnit.SECONDS).join()); + assertTrue(listener.received.isEmpty()); + assertEquals(0, routeInvocations.get()); + } + + @Test + public void validatesHandshakeWhileSecuredConsumerIsStopped() throws Exception { + context.getRouteController().startRoute("secureWs"); context.getRouteController().stopRoute("secureWs"); - // the handshake completes without OAuth validation because no consumer is set on the shared handler - CompletableFuture<Integer> closeCode = new CompletableFuture<>(); - WebSocket webSocket = connectWithoutAuthorization(closeCode); + CompletionException thrown = assertThrows(CompletionException.class, () -> connect(null, new RecordingListener())); + WebSocketHandshakeException handshake = assertInstanceOf(WebSocketHandshakeException.class, thrown.getCause()); + assertEquals(401, handshake.getResponse().statusCode()); - // once the secured consumer is back, the unauthenticated channel must not reach the route + WebSocket webSocket = connect("valid-token", new RecordingListener()); context.getRouteController().startRoute("secureWs"); + getMockEndpoint("mock:message").expectedBodiesReceived("hello"); webSocket.sendText("hello", true).orTimeout(5, TimeUnit.SECONDS).join(); - assertEquals(CloseMessage.MSG_VIOLATES_POLICY, closeCode.orTimeout(10, TimeUnit.SECONDS).join()); - assertEquals(0, routeInvocations.get()); + getMockEndpoint("mock:message").assertIsSatisfied(); } - private WebSocket connectWithoutAuthorization(CompletableFuture<Integer> closeCode) { - return HttpClient.newHttpClient().newWebSocketBuilder() - .buildAsync(URI.create("ws://localhost:" + getPort() + "/window"), new WebSocket.Listener() { - @Override - public void onOpen(WebSocket webSocket) { - webSocket.request(1); - } - - @Override - public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, String reason) { - closeCode.complete(statusCode); - return null; - } - }) + private WebSocket connect(String token, RecordingListener listener) { + WebSocket.Builder builder = HttpClient.newHttpClient().newWebSocketBuilder(); + if (token != null) { + builder.header("Authorization", "Bearer " + token); + } + return builder.buildAsync(URI.create("ws://localhost:" + getPort() + "/window"), listener) .orTimeout(5, TimeUnit.SECONDS).join(); } @@ -94,9 +112,34 @@ public class UndertowWsOAuthProfileConsumerWindowTest extends BaseUndertowTest { .to("undertow:ws://localhost:{{port}}/window?sendToAll=true"); from("undertow:ws://localhost:{{port}}/window?oauthProfile=myprofile").routeId("secureWs") + .autoStartup(false) .process(exchange -> routeInvocations.incrementAndGet()) .to("mock:message"); } }; } + + private static final class RecordingListener implements WebSocket.Listener { + + private final List<String> received = new CopyOnWriteArrayList<>(); + private final CompletableFuture<Integer> closeCode = new CompletableFuture<>(); + + @Override + public void onOpen(WebSocket webSocket) { + webSocket.request(1); + } + + @Override + public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last) { + received.add(data.toString()); + webSocket.request(1); + return null; + } + + @Override + public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, String reason) { + closeCode.complete(statusCode); + return null; + } + } } diff --git a/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsSecurityWithoutProviderTest.java b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsSecurityWithoutProviderTest.java new file mode 100644 index 000000000000..8da47b15535f --- /dev/null +++ b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsSecurityWithoutProviderTest.java @@ -0,0 +1,197 @@ +/* + * 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.camel.component.undertow.ws; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.WebSocket; +import java.net.http.WebSocketHandshakeException; +import java.nio.charset.StandardCharsets; +import java.util.Base64; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CompletionStage; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; + +import io.undertow.util.StatusCodes; +import io.undertow.websockets.core.CloseMessage; +import org.apache.camel.CamelContext; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.component.mock.MockEndpoint; +import org.apache.camel.component.undertow.BaseUndertowTest; +import org.apache.camel.component.undertow.StubOAuthTokenValidationFactory; +import org.apache.camel.component.undertow.UndertowBasicAuthHandler; +import org.apache.camel.spi.OAuthTokenValidationFactory; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * WebSocket endpoints apply the allowedRoles and handlers options without a security provider, as HTTP endpoints do. + */ +class UndertowWsSecurityWithoutProviderTest extends BaseUndertowTest { + + private final AtomicInteger lateRouteInvocations = new AtomicInteger(); + private final AtomicInteger lateBasicRouteInvocations = new AtomicInteger(); + + @Override + protected CamelContext createCamelContext() throws Exception { + CamelContext context = super.createCamelContext(); + context.getRegistry().bind(OAuthTokenValidationFactory.FACTORY, new StubOAuthTokenValidationFactory()); + context.getRegistry().bind("basicAuth", new UndertowBasicAuthHandler()); + context.getRegistry().bind("lateBasicAuth", new UndertowBasicAuthHandler()); + return context; + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("undertow:ws://localhost:{{port}}/roles?allowedRoles=user").to("mock:roles"); + + from("undertow:ws://localhost:{{port}}/oauthRoles?oauthProfile=myprofile&allowedRoles=user") + .to("mock:oauthRoles"); + + from("undertow:ws://localhost:{{port}}/basic?handlers=#basicAuth").to("mock:basic"); + + from("direct:late").to("undertow:ws://localhost:{{port}}/late?sendToAll=true"); + from("undertow:ws://localhost:{{port}}/late?allowedRoles=user").routeId("late").autoStartup(false) + .process(exchange -> lateRouteInvocations.incrementAndGet()); + + from("direct:lateBasic").to("undertow:ws://localhost:{{port}}/lateBasic?sendToAll=true"); + from("undertow:ws://localhost:{{port}}/lateBasic?handlers=#lateBasicAuth").routeId("lateBasic") + .autoStartup(false) + .process(exchange -> lateBasicRouteInvocations.incrementAndGet()); + } + }; + } + + @Test + void allowedRolesWithoutProviderRefusesTheUpgrade() { + assertRefused("/roles", null, StatusCodes.FORBIDDEN); + } + + @Test + void allowedRolesWithoutProviderRefusesAValidBearerToken() { + assertRefused("/oauthRoles", "Bearer valid-token", StatusCodes.FORBIDDEN); + } + + @Test + void handlersRunBeforeTheUpgrade() throws Exception { + assertRefused("/basic", null, StatusCodes.UNAUTHORIZED); + + getMockEndpoint("mock:basic").expectedBodiesReceived("hello"); + String credentials = Base64.getEncoder().encodeToString("guest:secret".getBytes(StandardCharsets.UTF_8)); + WebSocket webSocket = connect("/basic", "Basic " + credentials, new RecordingListener()); + webSocket.sendText("hello", true).join(); + + MockEndpoint.assertIsSatisfied(context); + webSocket.sendClose(WebSocket.NORMAL_CLOSURE, "done").join(); + } + + @Test + void connectionOpenedBeforeTheConsumerStartedIsNotServed() throws Exception { + // only the producer uses the path, and it has no security settings, so the upgrade is not checked + RecordingListener listener = new RecordingListener(); + WebSocket webSocket = connect("/late", null, listener); + + context.getRouteController().startRoute("late"); + + // the connection does not get what the producer sends + template.sendBody("direct:late", "broadcast"); + // and its messages do not reach the route: it is closed instead + webSocket.sendText("hello", true).join(); + + assertEquals(CloseMessage.MSG_VIOLATES_POLICY, listener.closeCode.orTimeout(10, TimeUnit.SECONDS).join()); + assertTrue(listener.received.isEmpty()); + assertEquals(0, lateRouteInvocations.get()); + } + + @Test + void connectionOpenedBeforeAConsumerWithHandlersStartedIsNotServed() throws Exception { + // only the producer uses the path, and it has no security settings, so the upgrade does not go through the + // handlers of the consumer + RecordingListener listener = new RecordingListener(); + WebSocket webSocket = connect("/lateBasic", null, listener); + + context.getRouteController().startRoute("lateBasic"); + + // the connection does not get what the producer sends + template.sendBody("direct:lateBasic", "broadcast"); + // and its messages do not reach the route: it is closed instead + webSocket.sendText("hello", true).join(); + + assertEquals(CloseMessage.MSG_VIOLATES_POLICY, listener.closeCode.orTimeout(10, TimeUnit.SECONDS).join()); + assertTrue(listener.received.isEmpty()); + assertEquals(0, lateBasicRouteInvocations.get()); + + // a connection that goes through the handlers is served + String credentials = Base64.getEncoder().encodeToString("guest:secret".getBytes(StandardCharsets.UTF_8)); + WebSocket authenticated = connect("/lateBasic", "Basic " + credentials, new RecordingListener()); + authenticated.sendText("hello", true).join(); + await().atMost(10, TimeUnit.SECONDS).until(() -> lateBasicRouteInvocations.get() == 1); + authenticated.sendClose(WebSocket.NORMAL_CLOSURE, "done").join(); + } + + private void assertRefused(String path, String authorization, int statusCode) { + CompletionException thrown + = assertThrows(CompletionException.class, () -> connect(path, authorization, new RecordingListener())); + WebSocketHandshakeException handshake = assertInstanceOf(WebSocketHandshakeException.class, thrown.getCause()); + assertEquals(statusCode, handshake.getResponse().statusCode()); + } + + private WebSocket connect(String path, String authorization, WebSocket.Listener listener) { + WebSocket.Builder builder = HttpClient.newHttpClient().newWebSocketBuilder(); + if (authorization != null) { + builder.header("Authorization", authorization); + } + return builder.buildAsync(URI.create("ws://localhost:" + getPort() + path), listener) + .orTimeout(5, TimeUnit.SECONDS).join(); + } + + private static final class RecordingListener implements WebSocket.Listener { + + private final List<String> received = new CopyOnWriteArrayList<>(); + private final CompletableFuture<Integer> closeCode = new CompletableFuture<>(); + + @Override + public void onOpen(WebSocket webSocket) { + webSocket.request(1); + } + + @Override + public CompletionStage<?> onText(WebSocket webSocket, CharSequence data, boolean last) { + received.add(data.toString()); + webSocket.request(1); + return null; + } + + @Override + public CompletionStage<?> onClose(WebSocket webSocket, int statusCode, String reason) { + closeCode.complete(statusCode); + return null; + } + } +} diff --git a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc index 01071fd0570e..3dab963a9edd 100644 --- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc +++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc @@ -29,6 +29,24 @@ binding mode allowed it and marshalled the body, so a verb producing `text/plain a binary body failed. When `produces` lists several media types, the first json (or xml) type is used, otherwise the first one. Wildcards such as `*/*` are skipped, falling back to `application/json` or `application/xml`. +=== camel-undertow - WebSocket endpoints apply the security provider, allowed roles and handlers + +WebSocket endpoints (`ws://` and `wss://`) now apply the `UndertowSecurityProvider` (the `securityConfiguration` or +`securityProvider` option), the `allowedRoles` option, and the `handlers` and `accessLog` options, as HTTP endpoints +do. Previously these options had no effect on WebSocket endpoints. They apply to the upgrade request: a client that +the security provider rejects, or that connects to an endpoint with `allowedRoles` but no security provider, now gets +the same status code as on an HTTP endpoint (for example `403`) and no connection is opened. The headers that the +security provider adds are set on the exchanges of the connection. + +All the endpoints of a WebSocket path share its connections. The settings of the consumer apply to the path, +including the messages sent by producers on that path, and keep applying while the consumer is stopped; a path that +only has producers applies their settings. This also applies to the `oauthProfile` option: while the consumer is +stopped, upgrade requests to its path are still validated. A connection that was opened before these settings +applied does not receive the messages of the producers, and it is closed when it sends a message. + +A WebSocket client that connected without satisfying the configured security provider or allowed roles is now +rejected: it must authenticate, or the option must be removed from the WebSocket endpoint. + == Upgrading from 4.22.0 to 4.22.1 === camel-hazelcast - ReplicatedHazelcastAggregationRepository now applies the default serialization filter
