From 35fd9c58098cca2fcf8720fbaba7cd72fe9d65e1 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Tue, 13 Apr 2021 17:36:01 -0400 Subject: [PATCH] GH-3533: Register WebSocket endpoints at runtime (#3548) * GH-3533: Register WebSocket endpoints at runtime Fixes https://github.com/spring-projects/spring-integration/issues/3533 * Rework `WebSocketIntegrationConfigurationInitializer` to register beans functional way to avoid reflection for Spring Native support * Move `IntegrationServletWebSocketHandlerRegistry` into a separate file for better readability * Implement `DestructionAwareBeanPostProcessor` for `IntegrationServletWebSocketHandlerRegistry` to track runtime bean registrations and removals * Introduce an `IntegrationDynamicWebSocketHandlerMapping` to manage runtime mapping registrations and removals * Add `servlet-api` dependency into `websocket` to be able to compile an `IntegrationDynamicWebSocketHandlerMapping` * Fix typo in the exception message of the `StandardIntegrationFlowRegistration` * Start dynamically added beans together with associated `IntegrationFlow` in the `StandardIntegrationFlowContext` * Document new feature * * Fix language in docs * Don't start those `SmartLifecycle`s together with a dynamic flow which are not `isAutoStartup()` * * Fix `TomcatWebSocketTestServer` to configure servlet for `loadOnStartup = 1` * Fix `WebSocketDslTests` to make `clientWebSocketContainer.setAutoStartup(true)` --- build.gradle | 1 + .../StandardIntegrationFlowContext.java | 11 +- .../StandardIntegrationFlowRegistration.java | 8 +- .../IntegrationWebSocketContainer.java | 4 + ...grationDynamicWebSocketHandlerMapping.java | 56 +++++++ ...rationServletWebSocketHandlerRegistry.java | 142 +++++++++++++++++ ...etIntegrationConfigurationInitializer.java | 87 +++++------ .../WebSocketInboundChannelAdapter.java | 8 +- .../websocket/TomcatWebSocketTestServer.java | 7 +- .../websocket/dsl/WebSocketDslTests.java | 147 ++++++++++++++++++ src/reference/asciidoc/web-sockets.adoc | 51 +++++- src/reference/asciidoc/whats-new.adoc | 7 + 12 files changed, 460 insertions(+), 69 deletions(-) create mode 100644 spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java create mode 100644 spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationServletWebSocketHandlerRegistry.java create mode 100644 spring-integration-websocket/src/test/java/org/springframework/integration/websocket/dsl/WebSocketDslTests.java diff --git a/build.gradle b/build.gradle index d40c10000d..844f64517b 100644 --- a/build.gradle +++ b/build.gradle @@ -825,6 +825,7 @@ project('spring-integration-websocket') { api project(':spring-integration-core') api 'org.springframework:spring-websocket' optionalApi 'org.springframework:spring-webmvc' + providedImplementation "javax.servlet:javax.servlet-api:$servletApiVersion" testImplementation project(':spring-integration-event') testImplementation "org.apache.tomcat.embed:tomcat-embed-websocket:$tomcatVersion" diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowContext.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowContext.java index d3759191f1..c760735f11 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowContext.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowContext.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2020 the original author or authors. + * Copyright 2016-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -31,6 +31,7 @@ import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.AbstractBeanDefinition; import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.context.SmartLifecycle; import org.springframework.core.type.MethodMetadata; import org.springframework.integration.core.MessagingTemplate; import org.springframework.integration.dsl.IntegrationFlow; @@ -121,9 +122,15 @@ public final class StandardIntegrationFlowContext implements IntegrationFlowCont final String theFlowId = flowId; builder.additionalBeans.forEach((key, value) -> registerBean(key, value, theFlowId)); - IntegrationFlowRegistration registration = new StandardIntegrationFlowRegistration(integrationFlow, this, flowId); + IntegrationFlowRegistration registration = + new StandardIntegrationFlowRegistration(integrationFlow, this, flowId); if (builder.autoStartup) { registration.start(); + builder.additionalBeans.keySet() + .stream() + .filter(SmartLifecycle.class::isInstance) + .filter((lifecycle) -> ((SmartLifecycle) lifecycle).isAutoStartup()) + .forEach((lifecycle) -> ((SmartLifecycle) lifecycle).start()); } this.registry.put(flowId, registration); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowRegistration.java b/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowRegistration.java index 12eb31e3ae..4c5e83784c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowRegistration.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/dsl/context/StandardIntegrationFlowRegistration.java @@ -1,5 +1,5 @@ /* - * Copyright 2018-2020 the original author or authors. + * Copyright 2018-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -52,7 +52,9 @@ class StandardIntegrationFlowRegistration implements IntegrationFlowRegistration private ConfigurableListableBeanFactory beanFactory; - StandardIntegrationFlowRegistration(IntegrationFlow integrationFlow, IntegrationFlowContext integrationFlowContext, String id) { + StandardIntegrationFlowRegistration(IntegrationFlow integrationFlow, IntegrationFlowContext integrationFlowContext, + String id) { + this.integrationFlow = integrationFlow; this.integrationFlowContext = integrationFlowContext; this.id = id; @@ -81,7 +83,7 @@ class StandardIntegrationFlowRegistration implements IntegrationFlowRegistration throw new IllegalStateException("Only 'IntegrationFlow' instances started from the 'MessageChannel' " + "(e.g. extracted from 'IntegrationFlow' Lambdas) can be used " + "for direct 'send' operation. " + - "But [" + this.integrationFlow + "] ins't one of them.\n" + + "But [" + this.integrationFlow + "] isn't one of them.\n" + "Consider 'BeanFactory.getBean()' usage for sending messages " + "to the required 'MessageChannel'."); } diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/IntegrationWebSocketContainer.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/IntegrationWebSocketContainer.java index 42865a95f1..806b9001f9 100644 --- a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/IntegrationWebSocketContainer.java +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/IntegrationWebSocketContainer.java @@ -104,6 +104,10 @@ public abstract class IntegrationWebSocketContainer implements DisposableBean { } } + public WebSocketHandler getWebSocketHandler() { + return this.webSocketHandler; + } + public List getSubProtocols() { List protocols = new ArrayList<>(); if (this.messageListener != null) { diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java new file mode 100644 index 0000000000..b7e018bfa0 --- /dev/null +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationDynamicWebSocketHandlerMapping.java @@ -0,0 +1,56 @@ +/* + * Copyright 2021 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://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.springframework.integration.websocket.config; + +import java.util.HashMap; +import java.util.Map; + +import javax.servlet.http.HttpServletRequest; + +import org.springframework.web.HttpRequestHandler; +import org.springframework.web.servlet.HandlerExecutionChain; +import org.springframework.web.servlet.handler.AbstractHandlerMapping; + +/** + * The {@link AbstractHandlerMapping} implementation for dynamic WebSocket endpoint registrations in Spring Integration. + *

+ * TODO until https://github.com/spring-projects/spring-framework/issues/26798 + * + * @author Artem Bilan + * + * @since 5.5 + */ +class IntegrationDynamicWebSocketHandlerMapping extends AbstractHandlerMapping { + + private final Map handlerMap = new HashMap<>(); + + @Override + protected Object getHandlerInternal(HttpServletRequest request) { + String lookupPath = initLookupPath(request); + HttpRequestHandler httpRequestHandler = this.handlerMap.get(lookupPath); + return httpRequestHandler != null ? new HandlerExecutionChain(httpRequestHandler) : null; + } + + void registerHandler(String path, HttpRequestHandler httpHandler) { + this.handlerMap.put(path, httpHandler); + } + + void unregisterHandler(String path) { + this.handlerMap.remove(path); + } + +} diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationServletWebSocketHandlerRegistry.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationServletWebSocketHandlerRegistry.java new file mode 100644 index 0000000000..898d9b72c2 --- /dev/null +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/IntegrationServletWebSocketHandlerRegistry.java @@ -0,0 +1,142 @@ +/* + * Copyright 2021 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://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.springframework.integration.websocket.config; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import org.springframework.beans.BeansException; +import org.springframework.beans.factory.config.DestructionAwareBeanPostProcessor; +import org.springframework.context.ApplicationContext; +import org.springframework.context.ApplicationContextAware; +import org.springframework.integration.websocket.ServerWebSocketContainer; +import org.springframework.scheduling.TaskScheduler; +import org.springframework.util.CollectionUtils; +import org.springframework.util.MultiValueMap; +import org.springframework.web.HttpRequestHandler; +import org.springframework.web.servlet.handler.AbstractHandlerMapping; +import org.springframework.web.socket.WebSocketHandler; +import org.springframework.web.socket.config.annotation.ServletWebSocketHandlerRegistration; +import org.springframework.web.socket.config.annotation.ServletWebSocketHandlerRegistry; +import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistration; + +/** + * The {@link ServletWebSocketHandlerRegistry} extension for Spring Integration purpose, especially + * a dynamic WebSocket endpoint registrations. + * + * @author Artem Bilan + * + * @since 5.5 + */ +class IntegrationServletWebSocketHandlerRegistry extends ServletWebSocketHandlerRegistry + implements ApplicationContextAware, DestructionAwareBeanPostProcessor { + + private final Map> dynamicRegistrations = new HashMap<>(); + + private ApplicationContext applicationContext; + + private volatile IntegrationDynamicWebSocketHandlerMapping dynamicHandlerMapping; + + IntegrationServletWebSocketHandlerRegistry() { + } + + @Override + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + this.applicationContext = applicationContext; + } + + + @Override + protected boolean requiresTaskScheduler() { // NOSONAR visibility + return super.requiresTaskScheduler(); + } + + @Override + protected void setTaskScheduler(TaskScheduler scheduler) { // NOSONAR visibility + super.setTaskScheduler(scheduler); + } + + @Override + public AbstractHandlerMapping getHandlerMapping() { + AbstractHandlerMapping originHandlerMapping = super.getHandlerMapping(); + originHandlerMapping.setApplicationContext(this.applicationContext); + this.dynamicHandlerMapping = this.applicationContext.getBean(IntegrationDynamicWebSocketHandlerMapping.class); + return originHandlerMapping; + } + + @Override + public WebSocketHandlerRegistration addHandler(WebSocketHandler handler, String... paths) { + if (this.dynamicHandlerMapping != null) { + IntegrationDynamicWebSocketHandlerRegistration registration = + new IntegrationDynamicWebSocketHandlerRegistration(); + registration.addHandler(handler, paths); + MultiValueMap mappings = registration.getMapping(); + for (Map.Entry> entry : mappings.entrySet()) { + HttpRequestHandler httpHandler = entry.getKey(); + List patterns = entry.getValue(); + this.dynamicRegistrations.put(handler, patterns); + for (String pattern : patterns) { + this.dynamicHandlerMapping.registerHandler(pattern, httpHandler); + } + } + return registration; + } + else { + return super.addHandler(handler, paths); + } + } + + @Override + public Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException { + if (this.dynamicHandlerMapping != null && bean instanceof ServerWebSocketContainer) { + ((ServerWebSocketContainer) bean).registerWebSocketHandlers(this); + } + return bean; + } + + @Override + public boolean requiresDestruction(Object bean) { + return bean instanceof ServerWebSocketContainer; + } + + @Override + public void postProcessBeforeDestruction(Object bean, String beanName) throws BeansException { + if (requiresDestruction(bean)) { + removeRegistration((ServerWebSocketContainer) bean); + } + } + + void removeRegistration(ServerWebSocketContainer serverWebSocketContainer) { + List patterns = this.dynamicRegistrations.remove(serverWebSocketContainer.getWebSocketHandler()); + if (this.dynamicHandlerMapping != null && !CollectionUtils.isEmpty(patterns)) { + for (String pattern : patterns) { + this.dynamicHandlerMapping.unregisterHandler(pattern); + } + } + } + + private static final class IntegrationDynamicWebSocketHandlerRegistration + extends ServletWebSocketHandlerRegistration { + + MultiValueMap getMapping() { + return getMappings(); + } + + } + +} diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/WebSocketIntegrationConfigurationInitializer.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/WebSocketIntegrationConfigurationInitializer.java index b3172bfba9..3ebe68ef32 100644 --- a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/WebSocketIntegrationConfigurationInitializer.java +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/config/WebSocketIntegrationConfigurationInitializer.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2019 the original author or authors. + * Copyright 2014-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -26,17 +26,14 @@ import org.springframework.beans.factory.config.AbstractFactoryBean; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.support.BeanDefinitionReaderUtils; import org.springframework.beans.factory.support.BeanDefinitionRegistry; -import org.springframework.context.ApplicationContext; -import org.springframework.context.ApplicationContextAware; +import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.integration.config.IntegrationConfigurationInitializer; -import org.springframework.scheduling.TaskScheduler; import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.util.ClassUtils; import org.springframework.web.servlet.HandlerMapping; -import org.springframework.web.servlet.handler.AbstractHandlerMapping; import org.springframework.web.socket.config.annotation.DelegatingWebSocketConfiguration; -import org.springframework.web.socket.config.annotation.ServletWebSocketHandlerRegistry; import org.springframework.web.socket.config.annotation.WebSocketConfigurer; /** @@ -60,7 +57,7 @@ public class WebSocketIntegrationConfigurationInitializer implements Integration @Override public void initialize(ConfigurableListableBeanFactory beanFactory) throws BeansException { if (beanFactory instanceof BeanDefinitionRegistry) { - this.registerEnableWebSocketIfNecessary((BeanDefinitionRegistry) beanFactory); + registerEnableWebSocketIfNecessary((BeanDefinitionRegistry) beanFactory); } else { LOGGER.warn("'DelegatingWebSocketConfiguration' isn't registered because 'beanFactory'" + @@ -87,22 +84,30 @@ public class WebSocketIntegrationConfigurationInitializer implements Integration private void registerEnableWebSocketIfNecessary(BeanDefinitionRegistry registry) { if (SERVLET_PRESENT) { if (!registry.containsBeanDefinition("defaultSockJsTaskScheduler")) { - BeanDefinitionBuilder sockJsTaskSchedulerBuilder = - BeanDefinitionBuilder.genericBeanDefinition(ThreadPoolTaskScheduler.class) - .addPropertyValue("threadNamePrefix", "SockJS-") - .addPropertyValue("poolSize", Runtime.getRuntime().availableProcessors()) - .addPropertyValue("removeOnCancelPolicy", true); - + ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler(); + taskScheduler.setThreadNamePrefix("SockJS-"); + taskScheduler.setPoolSize(Runtime.getRuntime().availableProcessors()); + taskScheduler.setRemoveOnCancelPolicy(true); registry.registerBeanDefinition("defaultSockJsTaskScheduler", - sockJsTaskSchedulerBuilder.getBeanDefinition()); + new RootBeanDefinition(ThreadPoolTaskScheduler.class, () -> taskScheduler)); } if (!registry.containsBeanDefinition(DelegatingWebSocketConfiguration.class.getName()) && !registry.containsBeanDefinition(WEB_SOCKET_HANDLER_MAPPING_BEAN_NAME)) { + + registry.registerBeanDefinition("integrationServletWebSocketHandlerRegistry", + new RootBeanDefinition(IntegrationServletWebSocketHandlerRegistry.class, + IntegrationServletWebSocketHandlerRegistry::new)); + + BeanDefinitionReaderUtils.registerWithGeneratedName( + new RootBeanDefinition(IntegrationDynamicWebSocketHandlerMapping.class, + IntegrationDynamicWebSocketHandlerMapping::new), + registry); + BeanDefinitionBuilder enableWebSocketBuilder = - BeanDefinitionBuilder.genericBeanDefinition(WebSocketHandlerMappingFactoryBean.class) - .setRole(BeanDefinition.ROLE_INFRASTRUCTURE) - .addPropertyReference("sockJsTaskScheduler", "defaultSockJsTaskScheduler"); + BeanDefinitionBuilder.genericBeanDefinition(WebSocketHandlerMappingFactoryBean.class, + () -> createWebSocketHandlerMapping((BeanFactory) registry)) + .setRole(BeanDefinition.ROLE_INFRASTRUCTURE); registry.registerBeanDefinition(WEB_SOCKET_HANDLER_MAPPING_BEAN_NAME, enableWebSocketBuilder.getBeanDefinition()); @@ -110,25 +115,22 @@ public class WebSocketIntegrationConfigurationInitializer implements Integration } } - private static class WebSocketHandlerMappingFactoryBean extends AbstractFactoryBean - implements ApplicationContextAware { + private static WebSocketHandlerMappingFactoryBean createWebSocketHandlerMapping(BeanFactory beanFactory) { + WebSocketHandlerMappingFactoryBean mappingFactoryBean = new WebSocketHandlerMappingFactoryBean(); + mappingFactoryBean.registry = + beanFactory.getBean("integrationServletWebSocketHandlerRegistry", + IntegrationServletWebSocketHandlerRegistry.class); + mappingFactoryBean.sockJsTaskScheduler = + beanFactory.getBean("defaultSockJsTaskScheduler", ThreadPoolTaskScheduler.class); + return mappingFactoryBean; + } - private final IntegrationServletWebSocketHandlerRegistry registry = - new IntegrationServletWebSocketHandlerRegistry(); + private static class WebSocketHandlerMappingFactoryBean extends AbstractFactoryBean { + + private IntegrationServletWebSocketHandlerRegistry registry; private ThreadPoolTaskScheduler sockJsTaskScheduler; - private ApplicationContext applicationContext; - - public void setSockJsTaskScheduler(ThreadPoolTaskScheduler sockJsTaskScheduler) { - this.sockJsTaskScheduler = sockJsTaskScheduler; - } - - @Override - public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { - this.applicationContext = applicationContext; - } - @Override protected HandlerMapping createInstance() { BeanFactory beanFactory = getBeanFactory(); @@ -141,9 +143,7 @@ public class WebSocketIntegrationConfigurationInitializer implements Integration if (this.registry.requiresTaskScheduler()) { this.registry.setTaskScheduler(this.sockJsTaskScheduler); } - AbstractHandlerMapping handlerMapping = this.registry.getHandlerMapping(); - handlerMapping.setApplicationContext(this.applicationContext); - return handlerMapping; + return this.registry.getHandlerMapping(); } @Override @@ -153,21 +153,4 @@ public class WebSocketIntegrationConfigurationInitializer implements Integration } - private static class IntegrationServletWebSocketHandlerRegistry extends ServletWebSocketHandlerRegistry { - - IntegrationServletWebSocketHandlerRegistry() { - } - - @Override - public boolean requiresTaskScheduler() { // NOSONAR visibility - return super.requiresTaskScheduler(); - } - - @Override - public void setTaskScheduler(TaskScheduler scheduler) { // NOSONAR visibility - super.setTaskScheduler(scheduler); - } - - } - } diff --git a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java index f195923463..a653351883 100644 --- a/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java +++ b/spring-integration-websocket/src/main/java/org/springframework/integration/websocket/inbound/WebSocketInboundChannelAdapter.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2020 the original author or authors. + * Copyright 2014-2021 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -280,7 +280,7 @@ public class WebSocketInboundChannelAdapter extends MessageProducerSupport public boolean isActive() { boolean active = super.isActive(); if (!active) { - logger.warn("MessageProducer '" + this + "' isn't started to accept WebSocket events."); + logger.warn(() -> "MessageProducer '" + this + "' isn't started to accept WebSocket events."); } return active; } @@ -308,8 +308,8 @@ public class WebSocketInboundChannelAdapter extends MessageProducerSupport if (this.useBroker) { this.brokerHandler.handleMessage(message); } - else if (logger.isDebugEnabled()) { - logger.debug("Messages with non 'SimpMessageType.MESSAGE' type are ignored for sending to the " + + else { + logger.debug(() -> "Messages with non 'SimpMessageType.MESSAGE' type are ignored for sending to the " + "'outputChannel'. They have to be emitted as 'ApplicationEvent's " + "from the 'SubProtocolHandler'. Or using 'AbstractBrokerMessageHandler'(useBroker = true) " + "from server side. Received message: " + message); diff --git a/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/TomcatWebSocketTestServer.java b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/TomcatWebSocketTestServer.java index aa5e9810c6..63ed6ba888 100644 --- a/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/TomcatWebSocketTestServer.java +++ b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/TomcatWebSocketTestServer.java @@ -20,6 +20,7 @@ import java.io.File; import java.io.IOException; import org.apache.catalina.Context; +import org.apache.catalina.Wrapper; import org.apache.catalina.startup.Tomcat; import org.apache.tomcat.websocket.server.WsContextListener; @@ -49,8 +50,10 @@ public class TomcatWebSocketTestServer implements InitializingBean, DisposableBe Context context = this.tomcatServer.addContext("", System.getProperty("java.io.tmpdir")); context.addApplicationListener(WsContextListener.class.getName()); - Tomcat.addServlet(context, "dispatcherServlet", new DispatcherServlet(this.serverContext)) - .setAsyncSupported(true); + Wrapper dispatcherServlet = + Tomcat.addServlet(context, "dispatcherServlet", new DispatcherServlet(this.serverContext)); + dispatcherServlet.setAsyncSupported(true); + dispatcherServlet.setLoadOnStartup(1); context.addServletMappingDecoded("/", "dispatcherServlet"); } diff --git a/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/dsl/WebSocketDslTests.java b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/dsl/WebSocketDslTests.java new file mode 100644 index 0000000000..052ffbdefb --- /dev/null +++ b/spring-integration-websocket/src/test/java/org/springframework/integration/websocket/dsl/WebSocketDslTests.java @@ -0,0 +1,147 @@ +/* + * Copyright 2021 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://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.springframework.integration.websocket.dsl; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; + +import javax.websocket.DeploymentException; + +import org.junit.jupiter.api.Test; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.dsl.IntegrationFlows; +import org.springframework.integration.dsl.context.IntegrationFlowContext; +import org.springframework.integration.websocket.ClientWebSocketContainer; +import org.springframework.integration.websocket.ServerWebSocketContainer; +import org.springframework.integration.websocket.TomcatWebSocketTestServer; +import org.springframework.integration.websocket.inbound.WebSocketInboundChannelAdapter; +import org.springframework.integration.websocket.outbound.WebSocketOutboundMessageHandler; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHandlingException; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; +import org.springframework.web.context.support.AnnotationConfigWebApplicationContext; +import org.springframework.web.socket.client.WebSocketClient; +import org.springframework.web.socket.client.standard.StandardWebSocketClient; +import org.springframework.web.socket.server.HandshakeHandler; +import org.springframework.web.socket.server.standard.TomcatRequestUpgradeStrategy; +import org.springframework.web.socket.server.support.DefaultHandshakeHandler; + +@SpringJUnitConfig(classes = WebSocketDslTests.ClientConfig.class) +@DirtiesContext +public class WebSocketDslTests { + + @Autowired + TomcatWebSocketTestServer server; + + @Autowired + WebSocketClient webSocketClient; + + @Autowired + IntegrationFlowContext integrationFlowContext; + + @Test + void testDynamicServerEndpointRegistration() { + // Dynamic server flow + AnnotationConfigWebApplicationContext serverContext = this.server.getServerContext(); + IntegrationFlowContext serverIntegrationFlowContext = serverContext.getBean(IntegrationFlowContext.class); + ServerWebSocketContainer serverWebSocketContainer = + new ServerWebSocketContainer("/dynamic") + .setHandshakeHandler(serverContext.getBean(HandshakeHandler.class)); + + WebSocketInboundChannelAdapter webSocketInboundChannelAdapter = + new WebSocketInboundChannelAdapter(serverWebSocketContainer); + + QueueChannel dynamicRequestsChannel = new QueueChannel(); + + IntegrationFlow serverFlow = + IntegrationFlows.from(webSocketInboundChannelAdapter) + .channel(dynamicRequestsChannel) + .get(); + + IntegrationFlowContext.IntegrationFlowRegistration dynamicServerFlow = + serverIntegrationFlowContext.registration(serverFlow) + .addBean(serverWebSocketContainer) + .register(); + + // Dynamic client flow + ClientWebSocketContainer clientWebSocketContainer = + new ClientWebSocketContainer(this.webSocketClient, this.server.getWsBaseUrl() + "/dynamic"); + clientWebSocketContainer.setAutoStartup(true); + WebSocketOutboundMessageHandler webSocketOutboundMessageHandler = + new WebSocketOutboundMessageHandler(clientWebSocketContainer); + + IntegrationFlow clientFlow = flow -> flow.handle(webSocketOutboundMessageHandler); + + IntegrationFlowContext.IntegrationFlowRegistration dynamicClientFlow = + this.integrationFlowContext.registration(clientFlow) + .addBean(clientWebSocketContainer) + .register(); + + dynamicClientFlow.getInputChannel().send(new GenericMessage<>("dynamic test")); + + Message result = dynamicRequestsChannel.receive(10_000); + assertThat(result).isNotNull() + .extracting(Message::getPayload) + .isEqualTo("dynamic test"); + + dynamicServerFlow.destroy(); + + assertThatExceptionOfType(MessageHandlingException.class) + .isThrownBy(() -> dynamicClientFlow.getInputChannel().send(new GenericMessage<>("another test"))) + .withCauseInstanceOf(DeploymentException.class) + .withMessageContaining("The HTTP response from the server [404]"); + + dynamicClientFlow.destroy(); + } + + @Configuration + @EnableIntegration + public static class ClientConfig { + + @Bean + public TomcatWebSocketTestServer server() { + return new TomcatWebSocketTestServer(WebSocketDslTests.ServerConfig.class); + } + + @Bean + public WebSocketClient webSocketClient() { + return new StandardWebSocketClient(); + } + + } + + + @Configuration + @EnableIntegration + static class ServerConfig { + + @Bean + public DefaultHandshakeHandler handshakeHandler() { + return new DefaultHandshakeHandler(new TomcatRequestUpgradeStrategy()); + } + + } + +} diff --git a/src/reference/asciidoc/web-sockets.adoc b/src/reference/asciidoc/web-sockets.adoc index 1f583ea85b..77369da48f 100644 --- a/src/reference/asciidoc/web-sockets.adoc +++ b/src/reference/asciidoc/web-sockets.adoc @@ -3,14 +3,14 @@ Starting with version 4.1, Spring Integration has WebSocket support. It is based on the architecture, infrastructure, and API from the Spring Framework's `web-socket` module. -Therefore, many of Spring WebSocket's components (such as `SubProtocolHandler` or `WebSocketClient`) and configuration options (such as `@EnableWebSocketMessageBroker`) can be reused within Spring Integration. +Therefore, many of Spring WebSocket's components (such as `SubProtocolHandler` or `WebSocketClient`) and configuration options (such as `@EnableWebSocketMessageBroker`) can be reused within Spring Integration. For more information, see the https://docs.spring.io/spring/docs/current/spring-framework-reference/web.html#websocket[Spring Framework WebSocket Support] chapter in the Spring Framework reference manual. You need to include this dependency into your project: ==== +[source, xml, subs="normal", role="primary"] .Maven -[source, xml, subs="normal"] ---- org.springframework.integration @@ -18,9 +18,8 @@ You need to include this dependency into your project: {project-version} ---- - +[source, groovy, subs="normal", role="secondary"] .Gradle -[source, groovy, subs="normal"] ---- compile "org.springframework.integration:spring-integration-websocket:{project-version}" ---- @@ -102,8 +101,7 @@ The adapter delegates to the `SubProtocolHandlerRegistry` to determine the appro NOTE: By default, the `WebSocketInboundChannelAdapter` relies only on the raw `PassThruSubProtocolHandler` implementation, which converts the `WebSocketMessage` to a `Message`. The `WebSocketInboundChannelAdapter` accepts and sends to the underlying integration flow only `Message` instances that have `SimpMessageType.MESSAGE` or an empty `simpMessageType` header. -All other `Message` types are handled through the `ApplicationEvent` instances emitted from a `SubProtocolHandler` implementation (such as -`StompSubProtocolHandler`). +All other `Message` types are handled through the `ApplicationEvent` instances emitted from a `SubProtocolHandler` implementation (such as `StompSubProtocolHandler`). On the server side, if the `@EnableWebSocketMessageBroker` configuration is present, you can configure `WebSocketInboundChannelAdapter` with the `useBroker = true` option. In this case, all `non-MESSAGE` `Message` types are delegated to the provided `AbstractBrokerMessageHandler`. @@ -391,3 +389,44 @@ For proper client side message preparation, you must inject an instance of the ` One problem with the default `StompSubProtocolHandler` is that it was designed for the server side, so it updates the `SEND` `stompCommand` header into `MESSAGE` (as required by the STOMP protocol for the server side). If the client does not send its messages in the proper `SEND` web socket frame, some STOMP brokers do not accept them. The purpose of the `ClientStompEncoder`, in this case, is to override the `stompCommand` header and set it to the `SEND` value before encoding the message to the `byte[]`. + +[[websocket-dynamic-endpoints]] +=== Dynamic WebSocket Endpoints Registration + +Starting with version 5.5, the WebSocket server endpoints (channel adapters based on a `ServerWebSocketContainer`) can now be registered (and removed) at runtime - the `paths` a `ServerWebSocketContainer` is mapped is exposed via `HandlerMapping` into a `DispatcherServlet` and accessible for WebSocket clients. +The <<./dsl.adoc#java-dsl-runtime-flows,Dynamic and Runtime Integration Flows>> support helps to register these endpoints in a transparent manner: + +==== +[source,java] +---- +@Autowired +IntegrationFlowContext integrationFlowContext; + +@Autowired +HandshakeHandler handshakeHandler; +... +ServerWebSocketContainer serverWebSocketContainer = + new ServerWebSocketContainer("/dynamic") + .setHandshakeHandler(this.handshakeHandler); + +WebSocketInboundChannelAdapter webSocketInboundChannelAdapter = + new WebSocketInboundChannelAdapter(serverWebSocketContainer); + +QueueChannel dynamicRequestsChannel = new QueueChannel(); + +IntegrationFlow serverFlow = + IntegrationFlows.from(webSocketInboundChannelAdapter) + .channel(dynamicRequestsChannel) + .get(); + +IntegrationFlowContext.IntegrationFlowRegistration dynamicServerFlow = + this.integrationFlowContext.registration(serverFlow) + .addBean(serverWebSocketContainer) + .register(); +... +dynamicServerFlow.destroy(); +---- +==== + +NOTE: It is important to call `.addBean(serverWebSocketContainer)` on the dynamic flow registration to add the instance of `ServerWebSocketContainer` into an `ApplicationContext` for endpoint registration. +When a dynamic flow registration is destroyed, the associated `ServerWebSocketContainer` instance is destroyed, too, as well as the respective endpoint registration, including URL path mappings. \ No newline at end of file diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 6921b4c70d..d0b0122d68 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -86,3 +86,10 @@ The `MongoDbMessageSourceSpec` was added into MongoDd Java DSL. An `update` option is now exposed on both the `MongoDbMessageSource` and `ReactiveMongoDbMessageSource` implementations. See <<./mongodb.adoc#mongodb,MongoDb Support>> for more information. + +[[x5.5-websocket]] +==== WebSockets Changes + +The WebSocket channel adapters based on `ServerWebSocketContainer` can now be registered and removed at runtime. + +See <<./web-sockets.adoc#web-sockets,WebSockets Support>> for more information.