This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch camel-4.18.x
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/camel-4.18.x by this push:
new f18591dc9c7e [backport camel-4.18.x] CAMEL-25375: camel-undertow -
Apply securityProvider, allowedRoles and handlers to WebSocket endpoints
(#27492)
f18591dc9c7e is described below
commit f18591dc9c7e4e36561bb7e8e3a2da655efdaed2
Author: Claus Ibsen <[email protected]>
AuthorDate: Wed Oct 7 21:30:27 2026 +0200
[backport camel-4.18.x] CAMEL-25375: camel-undertow - Apply
securityProvider, allowedRoles and handlers to WebSocket endpoints (#27492)
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
---
.../security/SpringSecurityWebSocketTest.java | 90 +++++++
.../src/main/docs/undertow-component.adoc | 20 +-
.../component/undertow/DefaultUndertowHost.java | 104 +++++---
.../component/undertow/UndertowComponent.java | 24 +-
.../camel/component/undertow/UndertowConsumer.java | 100 +++++---
.../camel/component/undertow/UndertowEndpoint.java | 40 +++
.../camel/component/undertow/UndertowHost.java | 19 ++
.../camel/component/undertow/UndertowProducer.java | 6 +-
.../undertow/handlers/CamelWebSocketHandler.java | 276 ++++++++++++++++++++-
.../CamelWebSocketHandlerSecuritySettingsTest.java | 87 +++++++
.../ProviderWithServletWebSocketEndpointTest.java | 91 +++++++
.../spi/ProviderWithServletWebSocketTest.java | 66 +++++
...ityProviderRolesFromComponentWebSocketTest.java | 78 ++++++
.../spi/SecurityProviderWebSocketTest.java | 160 ++++++++++++
.../spi/SecurityProviderWebSocketWrapTest.java | 89 +++++++
.../ws/UndertowWsSecurityWithoutProviderTest.java | 206 +++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_18.adoc | 18 ++
17 files changed, 1394 insertions(+), 80 deletions(-)
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 7c9daea911e7..a2799078795f 100644
--- a/components/camel-undertow/src/main/docs/undertow-component.adoc
+++ b/components/camel-undertow/src/main/docs/undertow-component.adoc
@@ -134,10 +134,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 5d629933397c..11b2b45c28cc 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) {
@@ -98,12 +109,14 @@ 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;
+ if (requiresServletContext(endpoint)) {
+ // deploy before the server starts, so that a failure
leaves no server running
+ deployServletContext();
}
+ undertow = builder.setHandler(exchange ->
entryHandler.handleRequest(exchange)).build();
LOG.info("Starting Undertow server on {}://{}:{}",
key.getSslContext() != null ? "https" : "http",
key.getHost(),
key.getPort());
@@ -123,10 +136,18 @@ public class DefaultUndertowHost implements UndertowHost {
// initialization again.
undertow.stop();
undertow = null;
+ if (deploymentManager != null) {
+ deploymentManager.undeploy();
+ deploymentManager = null;
+ }
throw e;
}
}
+ if (deploymentManager == null && requiresServletContext(endpoint))
{
+ // a later endpoint needs a servlet context: wrap the handler
of the running server
+ deployServletContext();
+ }
if (consumer != null && consumer.isRest()) {
restHandler.addConsumer(consumer);
return restHandler;
@@ -139,36 +160,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
@@ -188,11 +217,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 8140a0099537..80ee328854d7 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;
@@ -38,7 +37,9 @@ import io.undertow.util.HttpString;
import io.undertow.util.Methods;
import io.undertow.util.MimeMappings;
import io.undertow.util.StatusCodes;
+import io.undertow.websockets.core.CloseMessage;
import io.undertow.websockets.core.WebSocketChannel;
+import io.undertow.websockets.core.WebSockets;
import io.undertow.websockets.spi.WebSocketHttpExchange;
import org.apache.camel.AsyncCallback;
import org.apache.camel.Exchange;
@@ -92,11 +93,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
@@ -111,26 +108,13 @@ 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);
- 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
@@ -138,6 +122,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;
@@ -193,19 +198,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;
}
@@ -292,10 +287,14 @@ public class UndertowConsumer extends DefaultConsumer
implements HttpHandler, Su
* @param message the message received via the {@link
WebSocketChannel}
*/
public void sendMessage(final String connectionKey, WebSocketChannel
channel, final Object message) {
+ if (rejectUnauthenticatedWebSocketChannel(connectionKey, channel)) {
+ return;
+ }
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);
@@ -317,9 +316,13 @@ public class UndertowConsumer extends DefaultConsumer
implements HttpHandler, Su
*/
public void sendEventNotification(
String connectionKey, WebSocketHttpExchange transportExchange,
WebSocketChannel channel, EventType eventType) {
+ if (rejectUnauthenticatedWebSocketChannel(connectionKey, channel)) {
+ return;
+ }
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);
@@ -371,4 +374,29 @@ public class UndertowConsumer extends DefaultConsumer
implements HttpHandler, Su
return exchange;
}
+ /**
+ * Fails closed for WebSocket channels whose handshake did not pass the
security checks of this consumer's endpoint
+ * (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 (CamelWebSocketHandler.isAuthenticated(channel, getEndpoint())) {
+ return false;
+ }
+ 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);
+ }
+ }
+
}
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 3cd6cc3ff99e..3143e8f22321 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;
@@ -478,6 +481,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 00feeac6a644..39e8b8e41739 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
@@ -23,13 +23,17 @@ import java.io.Reader;
import java.io.StringReader;
import java.nio.ByteBuffer;
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;
@@ -38,11 +42,15 @@ 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.StatusCodes;
import io.undertow.websockets.WebSocketConnectionCallback;
import io.undertow.websockets.WebSocketProtocolHandshakeHandler;
import io.undertow.websockets.core.AbstractReceiveListener;
import io.undertow.websockets.core.BufferedBinaryMessage;
import io.undertow.websockets.core.BufferedTextMessage;
+import io.undertow.websockets.core.CloseMessage;
import io.undertow.websockets.core.WebSocketChannel;
import io.undertow.websockets.core.WebSockets;
import io.undertow.websockets.spi.WebSocketHttpExchange;
@@ -53,7 +61,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.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -66,11 +76,24 @@ import org.xnio.Pooled;
*/
public class CamelWebSocketHandler implements HttpHandler {
private static final Logger LOG =
LoggerFactory.getLogger(CamelWebSocketHandler.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;
@@ -79,6 +102,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();
@@ -130,9 +164,131 @@ public class CamelWebSocketHandler implements HttpHandler
{
*/
@Override
public void handleRequest(HttpServerExchange exchange) throws Exception {
+ 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 {
+ PathSecurity path = pathSecurity();
+ if (path.endpoints.stream().noneMatch(endpoint ->
requiresHandshakeResult(endpoint, path.consumer))) {
+ 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 : path.endpoints) {
+ 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);
+ }
+ }
+ if (requiresHandshakeResult(endpoint, path.consumer) &&
(!path.consumer || 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);
}
+ /**
+ * 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 PathSecurity pathSecurity() {
+ consumerLock.lock();
+ try {
+ if (consumerEndpoint != null) {
+ return new PathSecurity(List.of(consumerEndpoint), true);
+ }
+ } finally {
+ consumerLock.unlock();
+ }
+ return new
PathSecurity(producerEndpoints.stream().distinct().toList(), false);
+ }
+
+ /**
+ * Whether a channel must have passed the security provider, the allowed
roles or, for the endpoint of a consumer,
+ * the custom handlers of the endpoint during its handshake. Only
consumers run their custom handlers.
+ */
+ private static boolean requiresHandshakeResult(UndertowEndpoint endpoint,
boolean consumer) {
+ return endpoint.requiresAuthentication() || consumer &&
endpoint.getHandlers() != null;
+ }
+
+ private static boolean isAuthenticated(WebSocketChannel channel,
PathSecurity path) {
+ for (UndertowEndpoint endpoint : path.endpoints) {
+ if (!isAuthenticated(channel, endpoint, path.consumer)) {
+ return false;
+ }
+ }
+ return true;
+ }
+
+ /**
+ * Whether the handshake of the given channel passed the security checks
of the given endpoint of a consumer: 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) {
+ return isAuthenticated(channel, endpoint, true);
+ }
+
+ private static boolean isAuthenticated(WebSocketChannel channel,
UndertowEndpoint endpoint, boolean consumer) {
+ if (requiresHandshakeResult(endpoint, consumer)) {
+ return channel != null
+ && channel.getAttribute(HANDSHAKE_RESULT) instanceof
HandshakeResult result
+ && result.endpointUris.contains(endpoint.getEndpointUri());
+ }
+ return true;
+ }
+
+ /**
+ * 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();
+ }
+
/**
* Send the given {@code message} to one or more channels selected using
the given {@code peerFilter} within the
* given {@code timeout} and report the outcome to the given {@code
camelExchange} and {@code camelCallback}.
@@ -150,8 +306,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
+ PathSecurity path = pathSecurity();
+ List<WebSocketChannel> targetPeers =
delegate.getPeerConnections().stream()
+ .filter(peer -> isAuthenticated(peer, path))
+ .filter(peerFilter)
+ .collect(Collectors.toList());
if (targetPeers.isEmpty()) {
camelCallback.done(true);
return true;
@@ -169,6 +329,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) {
@@ -177,11 +346,63 @@ 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();
+ // endpoints that configure the same security configuration each get a
provider of their own, which are alike
+ boolean sameProvider = producerEndpoint.getSecurityProvider() ==
consumerEndpoint.getSecurityProvider()
+ || producerEndpoint.getSecurityConfiguration() != null
+ && producerEndpoint.getSecurityConfiguration() ==
consumerEndpoint.getSecurityConfiguration();
+ return ownSettings && (!sameProvider
+ || !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();
@@ -364,6 +585,46 @@ public class CamelWebSocketHandler implements HttpHandler {
}
+ /**
+ * The endpoints whose security settings apply to a path, and whether that
is the endpoint of its consumer.
+ */
+ private static final class PathSecurity {
+ private final List<UndertowEndpoint> endpoints;
+ private final boolean consumer;
+
+ private PathSecurity(List<UndertowEndpoint> endpoints, boolean
consumer) {
+ this.endpoints = endpoints;
+ this.consumer = consumer;
+ }
+ }
+
+ /**
+ * 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.
*/
@@ -375,6 +636,17 @@ public class CamelWebSocketHandler implements HttpHandler {
@Override
public void onConnect(WebSocketHttpExchange exchange, WebSocketChannel
channel) {
LOG.trace("onConnect {}", exchange);
+ HandshakeResult handshakeResult =
exchange.getAttachment(HANDSHAKE_RESULT_ATTACHMENT);
+ if (handshakeResult != null) {
+ channel.setAttribute(HANDSHAKE_RESULT, handshakeResult);
+ }
+ if (!isAuthenticated(channel, pathSecurity())) {
+ // 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);
channel.getReceiveSetter().set(receiveListener);
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..69c6187ea203
--- /dev/null
+++
b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/handlers/CamelWebSocketHandlerSecuritySettingsTest.java
@@ -0,0 +1,87 @@
+/*
+ * 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.component.undertow.spi.AbstractSecurityProviderTest;
+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"));
+ }
+
+ @Test
+ void producerWithTheSameSecurityConfigurationAsTheConsumer() throws
Exception {
+ // each endpoint that configures a security configuration gets a
provider instance of its own
+ Object configuration = new Object();
+ try (CamelContext context = new DefaultCamelContext()) {
+ context.start();
+ UndertowEndpoint consumerEndpoint = endpoint(context,
"allowedRoles=user", configuration);
+ UndertowEndpoint producerEndpoint = endpoint(context,
"allowedRoles=user&sendToAll=true", configuration);
+
assertFalse(CamelWebSocketHandler.hasUnusedSecuritySettings(consumerEndpoint,
producerEndpoint));
+
+ UndertowEndpoint otherProducerEndpoint = endpoint(context,
"allowedRoles=user&sendToAll=false", new Object());
+
assertTrue(CamelWebSocketHandler.hasUnusedSecuritySettings(consumerEndpoint,
otherProducerEndpoint));
+ }
+ }
+
+ private static UndertowEndpoint endpoint(CamelContext context, String
options, Object securityConfiguration) {
+ UndertowEndpoint endpoint
+ = context.getEndpoint("undertow:ws://localhost:8080/path?" +
options, UndertowEndpoint.class);
+ endpoint.setSecurityConfiguration(securityConfiguration);
+ endpoint.setSecurityProvider(new
AbstractSecurityProviderTest.MockSecurityProvider());
+ return endpoint;
+ }
+
+ 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/UndertowWsSecurityWithoutProviderTest.java
b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsSecurityWithoutProviderTest.java
new file mode 100644
index 000000000000..9f7ee35ec259
--- /dev/null
+++
b/components/camel-undertow/src/test/java/org/apache/camel/component/undertow/ws/UndertowWsSecurityWithoutProviderTest.java
@@ -0,0 +1,206 @@
+/*
+ * 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.UndertowBasicAuthHandler;
+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;
+
+/**
+ * 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("basicAuth", new
UndertowBasicAuthHandler());
+ context.getRegistry().bind("lateBasicAuth", new
UndertowBasicAuthHandler());
+ context.getRegistry().bind("producerBasicAuth", 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}}/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());
+
+ // handlers is a consumer option: a producer does not run it
+ from("direct:producerHandlers")
+
.to("undertow:ws://localhost:{{port}}/producerHandlers?handlers=#producerBasicAuth&sendToAll=true");
+
+
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 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();
+ }
+
+ @Test
+ void handlersOfAProducerDoNotGuardItsPath() {
+ RecordingListener listener = new RecordingListener();
+ WebSocket webSocket = connect("/producerHandlers", null, listener);
+
+ // the server registers the connection shortly after the client has
completed the upgrade
+ await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
+ template.sendBody("direct:producerHandlers", "update");
+ assertFalse(listener.received.isEmpty());
+ });
+ assertEquals("update", listener.received.get(0));
+ webSocket.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_18.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_18.adoc
index 86e21b7a320a..11d62d937428 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_18.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_18.adoc
@@ -89,6 +89,24 @@ level by default.
Note that JGroups deserializes the payload inside its own receive path, so
this check is a
defense-in-depth allow-list on the resulting body type. The JVM-wide
`jdk.serialFilter`, together
with a channel secured with `AUTH` and encryption, remain the primary
mitigations.
+
+=== 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. 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.18.3 to 4.18.4
=== camel-core - Multicast EIP honors UseOriginalAggregationStrategy