diff --git a/bus/pom.xml b/bus/pom.xml index 5c8eccfbcd..300945511c 100644 --- a/bus/pom.xml +++ b/bus/pom.xml @@ -23,6 +23,7 @@ spring-cloud-starter-bus-amqp spring-cloud-starter-bus-kafka spring-cloud-starter-bus-stream + spring-cloud-starter-bus-jgroups @@ -77,6 +78,11 @@ spring-cloud-stream ${project.version} + + org.jgroups + jgroups + 5.5.7.Final + diff --git a/bus/spring-cloud-starter-bus-jgroups/pom.xml b/bus/spring-cloud-starter-bus-jgroups/pom.xml new file mode 100644 index 0000000000..2bdf6f133f --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/pom.xml @@ -0,0 +1,48 @@ + + + 4.0.0 + + + org.springframework.cloud + spring-cloud-bus-parent + 5.1.0-SNAPSHOT + .. + + + spring-cloud-starter-bus-jgroups + spring-cloud-starter-bus-jgroups + Spring Cloud Starter + https://projects.spring.io/spring-cloud + + + Pivotal Software, Inc. + https://www.spring.io + + + + ${basedir}/../.. + + + + + org.jgroups + jgroups + + + tools.jackson.core + jackson-databind + + + org.springframework.cloud + spring-cloud-bus + + + org.springframework.boot + spring-boot-starter-test + test + + + + diff --git a/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusAutoConfiguration.java b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusAutoConfiguration.java new file mode 100644 index 0000000000..ae316471a2 --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusAutoConfiguration.java @@ -0,0 +1,55 @@ +/* + * Copyright 2015-present 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.cloud.bus.jgroups; + +import java.util.function.Consumer; + +import org.jgroups.JChannel; + +import tools.jackson.databind.ObjectMapper; + +import org.springframework.boot.autoconfigure.AutoConfiguration; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.cloud.bus.BusAutoConfiguration; +import org.springframework.cloud.bus.BusBridge; +import org.springframework.cloud.bus.BusStreamAutoConfiguration; +import org.springframework.cloud.bus.ConditionalOnBusEnabled; +import org.springframework.cloud.bus.event.RemoteApplicationEvent; +import org.springframework.context.annotation.Bean; + +@AutoConfiguration(before = { BusStreamAutoConfiguration.class, BusAutoConfiguration.class }) +@ConditionalOnBusEnabled +@ConditionalOnClass({ JChannel.class, ObjectMapper.class }) +@EnableConfigurationProperties(JGroupsBusProperties.class) +public class JGroupsBusAutoConfiguration { + + @Bean + @ConditionalOnMissingBean + JGroupsChannelFactory jGroupsChannelFactory() { + return new JGroupsChannelFactory(); + } + + @Bean + @ConditionalOnMissingBean(BusBridge.class) + JGroupsBusBridge jGroupsBusBridge(JGroupsBusProperties properties, ObjectMapper objectMapper, + Consumer eventConsumer, JGroupsChannelFactory channelFactory) { + return new JGroupsBusBridge(properties, objectMapper, eventConsumer, channelFactory); + } + +} diff --git a/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusBridge.java b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusBridge.java new file mode 100644 index 0000000000..76fd775453 --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusBridge.java @@ -0,0 +1,80 @@ +/* + * Copyright 2015-present 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.cloud.bus.jgroups; + +import java.nio.charset.StandardCharsets; +import java.util.function.Consumer; + +import org.jgroups.BytesMessage; +import org.jgroups.JChannel; +import org.jgroups.Message; +import org.jgroups.Receiver; + +import tools.jackson.databind.ObjectMapper; + +import org.springframework.cloud.bus.BusBridge; +import org.springframework.cloud.bus.event.RemoteApplicationEvent; + +public class JGroupsBusBridge implements BusBridge { + + private final JChannel channel; + + private final ObjectMapper objectMapper; + + public JGroupsBusBridge(JGroupsBusProperties properties, ObjectMapper objectMapper, + Consumer eventConsumer, JGroupsChannelFactory channelFactory) { + try { + this.objectMapper = objectMapper; + this.channel = channelFactory.create(); + this.channel.setReceiver(new Receiver() { + @Override + public void receive(Message message) { + try { + String payload = new String(message.getArray(), message.getOffset(), message.getLength(), + StandardCharsets.UTF_8); + RemoteApplicationEvent event = JGroupsBusBridge.this.objectMapper.readValue(payload, + RemoteApplicationEvent.class); + eventConsumer.accept(event); + } + catch (Exception ex) { + throw new IllegalStateException("Unable to receive event from JGroups", ex); + } + } + }); + this.channel.connect(properties.getClusterName()); + } + catch (Exception ex) { + throw new IllegalStateException("Unable to initialize JGroups channel", ex); + } + } + + @Override + public void send(RemoteApplicationEvent event) { + try { + byte[] payload = this.objectMapper.writeValueAsBytes(event); + this.channel.send(new BytesMessage(null, payload)); + } + catch (Exception ex) { + throw new IllegalStateException("Unable to send event over JGroups", ex); + } + } + + public void close() { + this.channel.close(); + } + +} diff --git a/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusProperties.java b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusProperties.java new file mode 100644 index 0000000000..1680d363bc --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsBusProperties.java @@ -0,0 +1,36 @@ +/* + * Copyright 2015-present 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.cloud.bus.jgroups; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +@ConfigurationProperties(JGroupsBusProperties.PREFIX) +public class JGroupsBusProperties { + + public static final String PREFIX = "spring.cloud.bus.jgroups"; + + private String clusterName = "spring-cloud-bus"; + + public String getClusterName() { + return this.clusterName; + } + + public void setClusterName(String clusterName) { + this.clusterName = clusterName; + } + +} diff --git a/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsChannelFactory.java b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsChannelFactory.java new file mode 100644 index 0000000000..9ae3a62893 --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/src/main/java/org/springframework/cloud/bus/jgroups/JGroupsChannelFactory.java @@ -0,0 +1,27 @@ +/* + * Copyright 2015-present 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.cloud.bus.jgroups; + +import org.jgroups.JChannel; + +class JGroupsChannelFactory { + + JChannel create() throws Exception { + return new JChannel(); + } + +} diff --git a/bus/spring-cloud-starter-bus-jgroups/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports b/bus/spring-cloud-starter-bus-jgroups/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports new file mode 100644 index 0000000000..191c448ead --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/src/main/resources/META-INF/spring/org.springframework.boot.autoconfigure.AutoConfiguration.imports @@ -0,0 +1 @@ +org.springframework.cloud.bus.jgroups.JGroupsBusAutoConfiguration diff --git a/bus/spring-cloud-starter-bus-jgroups/src/test/java/org/springframework/cloud/bus/jgroups/JGroupsBusAutoConfigurationTests.java b/bus/spring-cloud-starter-bus-jgroups/src/test/java/org/springframework/cloud/bus/jgroups/JGroupsBusAutoConfigurationTests.java new file mode 100644 index 0000000000..af9e050e4f --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/src/test/java/org/springframework/cloud/bus/jgroups/JGroupsBusAutoConfigurationTests.java @@ -0,0 +1,64 @@ +/* + * Copyright 2015-present 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.cloud.bus.jgroups; + +import java.util.function.Consumer; + +import org.jgroups.JChannel; +import org.junit.jupiter.api.Test; + +import tools.jackson.databind.ObjectMapper; +import tools.jackson.databind.json.JsonMapper; + +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.cloud.bus.BusBridge; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.mock; + +class JGroupsBusAutoConfigurationTests { + + private final ApplicationContextRunner contextRunner = new ApplicationContextRunner() + .withPropertyValues("spring.cloud.bus.enabled=true") + .withUserConfiguration(JGroupsBusAutoConfiguration.class) + .withBean(ObjectMapper.class, () -> JsonMapper.builder().build()) + .withBean(Consumer.class, () -> mock(Consumer.class)) + .withBean(JGroupsChannelFactory.class, () -> { + JGroupsChannelFactory factory = mock(JGroupsChannelFactory.class); + JChannel channel = mock(JChannel.class); + try { + doReturn(channel).when(factory).create(); + } + catch (Exception ex) { + throw new IllegalStateException(ex); + } + return factory; + }); + + @Test + void createsJGroupsBusBridge() { + this.contextRunner.run(context -> assertThat(context).hasSingleBean(JGroupsBusBridge.class)); + } + + @Test + void backsOffWhenBusBridgeAlreadyExists() { + this.contextRunner.withBean(BusBridge.class, () -> mock(BusBridge.class)) + .run(context -> assertThat(context).doesNotHaveBean(JGroupsBusBridge.class)); + } + +} diff --git a/bus/spring-cloud-starter-bus-jgroups/src/test/java/org/springframework/cloud/bus/jgroups/JGroupsBusBridgeTests.java b/bus/spring-cloud-starter-bus-jgroups/src/test/java/org/springframework/cloud/bus/jgroups/JGroupsBusBridgeTests.java new file mode 100644 index 0000000000..c62c048870 --- /dev/null +++ b/bus/spring-cloud-starter-bus-jgroups/src/test/java/org/springframework/cloud/bus/jgroups/JGroupsBusBridgeTests.java @@ -0,0 +1,119 @@ +/* + * Copyright 2015-present 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.cloud.bus.jgroups; + +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; + +import org.jgroups.BytesMessage; +import org.jgroups.JChannel; +import org.jgroups.Receiver; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import tools.jackson.databind.ObjectMapper; +import tools.jackson.databind.json.JsonMapper; + +import org.springframework.cloud.bus.event.EnvironmentChangeRemoteApplicationEvent; +import org.springframework.cloud.bus.event.RemoteApplicationEvent; +import org.springframework.cloud.bus.jackson.SubtypeModule; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class JGroupsBusBridgeTests { + + private final JGroupsBusProperties properties = new JGroupsBusProperties(); + + private final ObjectMapper objectMapper = JsonMapper.builder() + .addModule(new SubtypeModule(EnvironmentChangeRemoteApplicationEvent.class)) + .build(); + + @Test + void sendsEventAsBytesMessage() throws Exception { + JChannel channel = mock(JChannel.class); + JGroupsChannelFactory channelFactory = mock(JGroupsChannelFactory.class); + when(channelFactory.create()).thenReturn(channel); + + JGroupsBusBridge bridge = new JGroupsBusBridge(this.properties, this.objectMapper, event -> { + }, channelFactory); + + RemoteApplicationEvent event = new EnvironmentChangeRemoteApplicationEvent("test", "test", (String) null, + Map.of()); + + bridge.send(event); + + verify(channel).send(any(BytesMessage.class)); + } + + @Test + void receivesEventAndDelegatesToConsumer() throws Exception { + JChannel channel = mock(JChannel.class); + JGroupsChannelFactory channelFactory = mock(JGroupsChannelFactory.class); + when(channelFactory.create()).thenReturn(channel); + + AtomicReference received = new AtomicReference<>(); + + new JGroupsBusBridge(this.properties, this.objectMapper, received::set, channelFactory); + + ArgumentCaptor receiver = ArgumentCaptor.forClass(Receiver.class); + verify(channel).setReceiver(receiver.capture()); + + EnvironmentChangeRemoteApplicationEvent event = new EnvironmentChangeRemoteApplicationEvent("test", "test", + (String) null, Map.of()); + + byte[] payload = this.objectMapper.writeValueAsBytes(event); + BytesMessage message = new BytesMessage(null, payload); + + receiver.getValue().receive(message); + + assertThat(received.get()).isNotNull(); + assertThat(received.get().getClass()).isEqualTo(EnvironmentChangeRemoteApplicationEvent.class); + } + + @Test + void connectsToConfiguredCluster() throws Exception { + JChannel channel = mock(JChannel.class); + JGroupsChannelFactory channelFactory = mock(JGroupsChannelFactory.class); + when(channelFactory.create()).thenReturn(channel); + + this.properties.setClusterName("test-cluster"); + + new JGroupsBusBridge(this.properties, this.objectMapper, event -> { + }, channelFactory); + + verify(channel).connect("test-cluster"); + } + + @Test + void closesChannel() throws Exception { + JChannel channel = mock(JChannel.class); + JGroupsChannelFactory channelFactory = mock(JGroupsChannelFactory.class); + when(channelFactory.create()).thenReturn(channel); + + JGroupsBusBridge bridge = new JGroupsBusBridge(this.properties, this.objectMapper, event -> { + }, channelFactory); + + bridge.close(); + + verify(channel).close(); + } + +}