Add Pulsar ConnectionDetails support
Add `ConnectionDetails` support for Apache Pulsar and provide adapters for Docker Compose and Testcontainers. See gh-37197
This commit is contained in:
@@ -277,4 +277,5 @@ tasks.named("checkSpringConfigurationMetadata").configure {
|
||||
|
||||
test {
|
||||
jvmArgs += "--add-opens=java.base/java.net=ALL-UNNAMED"
|
||||
jvmArgs += "--add-opens=java.base/sun.net=ALL-UNNAMED"
|
||||
}
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
/*
|
||||
* Copyright 2012-2023 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.boot.autoconfigure.pulsar;
|
||||
|
||||
/**
|
||||
* Adapts {@link PulsarProperties} to {@link PulsarConnectionDetails}.
|
||||
*
|
||||
* @author Chris Bono
|
||||
*/
|
||||
class PropertiesPulsarConnectionDetails implements PulsarConnectionDetails {
|
||||
|
||||
private final PulsarProperties pulsarProperties;
|
||||
|
||||
PropertiesPulsarConnectionDetails(PulsarProperties pulsarProperties) {
|
||||
this.pulsarProperties = pulsarProperties;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getPulsarBrokerUrl() {
|
||||
return this.pulsarProperties.getClient().getServiceUrl();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String getPulsarAdminUrl() {
|
||||
return this.pulsarProperties.getAdmin().getServiceUrl();
|
||||
}
|
||||
|
||||
}
|
||||
@@ -71,17 +71,31 @@ class PulsarConfiguration {
|
||||
this.propertiesMapper = new PulsarPropertiesMapper(properties);
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(PulsarConnectionDetails.class)
|
||||
PropertiesPulsarConnectionDetails pulsarConnectionDetails() {
|
||||
return new PropertiesPulsarConnectionDetails(this.properties);
|
||||
}
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(PulsarClientFactory.class)
|
||||
DefaultPulsarClientFactory pulsarClientFactory(ObjectProvider<PulsarClientBuilderCustomizer> customizersProvider) {
|
||||
DefaultPulsarClientFactory pulsarClientFactory(PulsarConnectionDetails connectionDetails,
|
||||
ObjectProvider<PulsarClientBuilderCustomizer> customizersProvider) {
|
||||
List<PulsarClientBuilderCustomizer> allCustomizers = new ArrayList<>();
|
||||
allCustomizers.add(this.propertiesMapper::customizeClientBuilder);
|
||||
allCustomizers.add((clientBuilder) -> this.applyConnectionDetails(connectionDetails, clientBuilder));
|
||||
allCustomizers.addAll(customizersProvider.orderedStream().toList());
|
||||
DefaultPulsarClientFactory clientFactory = new DefaultPulsarClientFactory(
|
||||
(clientBuilder) -> applyClientBuilderCustomizers(allCustomizers, clientBuilder));
|
||||
return clientFactory;
|
||||
}
|
||||
|
||||
private void applyConnectionDetails(PulsarConnectionDetails connectionDetails, ClientBuilder clientBuilder) {
|
||||
if (connectionDetails.getPulsarBrokerUrl() != null) {
|
||||
clientBuilder.serviceUrl(connectionDetails.getPulsarBrokerUrl());
|
||||
}
|
||||
}
|
||||
|
||||
private void applyClientBuilderCustomizers(List<PulsarClientBuilderCustomizer> customizers,
|
||||
ClientBuilder clientBuilder) {
|
||||
customizers.forEach((customizer) -> customizer.customize(clientBuilder));
|
||||
@@ -95,14 +109,21 @@ class PulsarConfiguration {
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean
|
||||
PulsarAdministration pulsarAdministration(
|
||||
PulsarAdministration pulsarAdministration(PulsarConnectionDetails connectionDetails,
|
||||
ObjectProvider<PulsarAdminBuilderCustomizer> pulsarAdminBuilderCustomizers) {
|
||||
List<PulsarAdminBuilderCustomizer> allCustomizers = new ArrayList<>();
|
||||
allCustomizers.add(this.propertiesMapper::customizeAdminBuilder);
|
||||
allCustomizers.add((adminBuilder) -> this.applyConnectionDetails(connectionDetails, adminBuilder));
|
||||
allCustomizers.addAll(pulsarAdminBuilderCustomizers.orderedStream().toList());
|
||||
return new PulsarAdministration((adminBuilder) -> applyAdminBuilderCustomizers(allCustomizers, adminBuilder));
|
||||
}
|
||||
|
||||
private void applyConnectionDetails(PulsarConnectionDetails connectionDetails, PulsarAdminBuilder adminBuilder) {
|
||||
if (connectionDetails.getPulsarAdminUrl() != null) {
|
||||
adminBuilder.serviceHttpUrl(connectionDetails.getPulsarAdminUrl());
|
||||
}
|
||||
}
|
||||
|
||||
private void applyAdminBuilderCustomizers(List<PulsarAdminBuilderCustomizer> customizers,
|
||||
PulsarAdminBuilder adminBuilder) {
|
||||
customizers.forEach((customizer) -> customizer.customize(adminBuilder));
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
/*
|
||||
* Copyright 2012-2023 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.boot.autoconfigure.pulsar;
|
||||
|
||||
import org.springframework.boot.autoconfigure.service.connection.ConnectionDetails;
|
||||
|
||||
/**
|
||||
* Details required to establish a connection to a Pulsar service.
|
||||
*
|
||||
* @author Chris Bono
|
||||
* @since 3.2.0
|
||||
*/
|
||||
public interface PulsarConnectionDetails extends ConnectionDetails {
|
||||
|
||||
/**
|
||||
* Returns the Pulsar service URL for the broker.
|
||||
* @return the Pulsar service URL for the broker
|
||||
*/
|
||||
String getPulsarBrokerUrl();
|
||||
|
||||
/**
|
||||
* Returns the Pulsar web URL for the admin endpoint.
|
||||
* @return the Pulsar web URL for the admin endpoint
|
||||
*/
|
||||
String getPulsarAdminUrl();
|
||||
|
||||
}
|
||||
@@ -53,6 +53,7 @@ final class PulsarPropertiesMapper {
|
||||
PulsarProperties.Client properties = this.properties.getClient();
|
||||
PropertyMapper map = PropertyMapper.get().alwaysApplyingWhenNonNull();
|
||||
map.from(properties::getServiceUrl).to(clientBuilder::serviceUrl);
|
||||
|
||||
map.from(properties::getConnectionTimeout).to(timeoutProperty(clientBuilder::connectionTimeout));
|
||||
map.from(properties::getOperationTimeout).to(timeoutProperty(clientBuilder::operationTimeout));
|
||||
map.from(properties::getLookupTimeout).to(timeoutProperty(clientBuilder::lookupTimeout));
|
||||
|
||||
@@ -0,0 +1,46 @@
|
||||
/*
|
||||
* Copyright 2012-2023 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.boot.autoconfigure.pulsar;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
|
||||
/**
|
||||
* Tests for {@link PropertiesPulsarConnectionDetails}.
|
||||
*
|
||||
* @author Chris Bono
|
||||
*/
|
||||
class PropertiesPulsarConnectionDetailsTests {
|
||||
|
||||
@Test
|
||||
void pulsarBrokerUrlIsObtainedFromPulsarProperties() {
|
||||
var pulsarProps = new PulsarProperties();
|
||||
pulsarProps.getClient().setServiceUrl("foo");
|
||||
var connectionDetails = new PropertiesPulsarConnectionDetails(pulsarProps);
|
||||
assertThat(connectionDetails.getPulsarBrokerUrl()).isEqualTo("foo");
|
||||
}
|
||||
|
||||
@Test
|
||||
void pulsarAdminUrlIsObtainedFromPulsarProperties() {
|
||||
var pulsarProps = new PulsarProperties();
|
||||
pulsarProps.getAdmin().setServiceUrl("foo");
|
||||
var connectionDetails = new PropertiesPulsarConnectionDetails(pulsarProps);
|
||||
assertThat(connectionDetails.getPulsarAdminUrl()).isEqualTo("foo");
|
||||
}
|
||||
|
||||
}
|
||||
@@ -114,6 +114,7 @@ class PulsarAutoConfigurationTests {
|
||||
@Test
|
||||
void autoConfiguresBeans() {
|
||||
this.contextRunner.run((context) -> assertThat(context).hasSingleBean(PulsarConfiguration.class)
|
||||
.hasSingleBean(PulsarConnectionDetails.class)
|
||||
.hasSingleBean(DefaultPulsarClientFactory.class)
|
||||
.hasSingleBean(PulsarClient.class)
|
||||
.hasSingleBean(PulsarAdministration.class)
|
||||
|
||||
@@ -51,6 +51,7 @@ import org.springframework.pulsar.function.PulsarFunctionAdministration;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.assertj.core.api.Assertions.entry;
|
||||
import static org.mockito.BDDMockito.given;
|
||||
import static org.mockito.Mockito.mock;
|
||||
|
||||
/**
|
||||
@@ -67,6 +68,15 @@ class PulsarConfigurationTests {
|
||||
.withConfiguration(AutoConfigurations.of(PulsarConfiguration.class))
|
||||
.withBean(PulsarClient.class, () -> mock(PulsarClient.class));
|
||||
|
||||
@Test
|
||||
void whenHasUserDefinedConnectionDetailsBeanDoesNotAutoConfigureBean() {
|
||||
PulsarConnectionDetails customConnectionDetails = mock(PulsarConnectionDetails.class);
|
||||
this.contextRunner
|
||||
.withBean("customPulsarConnectionDetails", PulsarConnectionDetails.class, () -> customConnectionDetails)
|
||||
.run((context) -> assertThat(context).getBean(PulsarConnectionDetails.class)
|
||||
.isSameAs(customConnectionDetails));
|
||||
}
|
||||
|
||||
@Nested
|
||||
class ClientTests {
|
||||
|
||||
@@ -86,17 +96,36 @@ class PulsarConfigurationTests {
|
||||
.run((context) -> assertThat(context).getBean(PulsarClient.class).isSameAs(customClient));
|
||||
}
|
||||
|
||||
@Test
|
||||
void whenConnectionDetailsAreNullTheyAreNotApplied() {
|
||||
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
|
||||
given(connectionDetails.getPulsarBrokerUrl()).willReturn(null);
|
||||
PulsarConfigurationTests.this.contextRunner.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
|
||||
.withPropertyValues("spring.pulsar.client.service-url=fromPropsCustomizer")
|
||||
.run((context) -> {
|
||||
DefaultPulsarClientFactory clientFactory = context.getBean(DefaultPulsarClientFactory.class);
|
||||
Customizers<PulsarClientBuilderCustomizer, ClientBuilder> customizers = Customizers
|
||||
.of(ClientBuilder.class, PulsarClientBuilderCustomizer::customize);
|
||||
assertThat(customizers.fromField(clientFactory, "customizer"))
|
||||
.callsInOrder(ClientBuilder::serviceUrl, "fromPropsCustomizer");
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
|
||||
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
|
||||
given(connectionDetails.getPulsarBrokerUrl()).willReturn("fromConnectionDetailsCustomizer");
|
||||
PulsarConfigurationTests.this.contextRunner
|
||||
.withUserConfiguration(PulsarClientBuilderCustomizersConfig.class)
|
||||
.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
|
||||
.withPropertyValues("spring.pulsar.client.service-url=fromPropsCustomizer")
|
||||
.run((context) -> {
|
||||
DefaultPulsarClientFactory clientFactory = context.getBean(DefaultPulsarClientFactory.class);
|
||||
Customizers<PulsarClientBuilderCustomizer, ClientBuilder> customizers = Customizers
|
||||
.of(ClientBuilder.class, PulsarClientBuilderCustomizer::customize);
|
||||
assertThat(customizers.fromField(clientFactory, "customizer")).callsInOrder(
|
||||
ClientBuilder::serviceUrl, "fromPropsCustomizer", "fromCustomizer1", "fromCustomizer2");
|
||||
ClientBuilder::serviceUrl, "fromPropsCustomizer", "fromConnectionDetailsCustomizer",
|
||||
"fromCustomizer1", "fromCustomizer2");
|
||||
});
|
||||
}
|
||||
|
||||
@@ -133,17 +162,35 @@ class PulsarConfigurationTests {
|
||||
.isSameAs(pulsarAdministration));
|
||||
}
|
||||
|
||||
@Test
|
||||
void whenConnectionDetailsAreNullTheyAreNotApplied() {
|
||||
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
|
||||
given(connectionDetails.getPulsarAdminUrl()).willReturn(null);
|
||||
PulsarConfigurationTests.this.contextRunner.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
|
||||
.withPropertyValues("spring.pulsar.admin.service-url=fromPropsCustomizer")
|
||||
.run((context) -> {
|
||||
PulsarAdministration pulsarAdmin = context.getBean(PulsarAdministration.class);
|
||||
Customizers<PulsarAdminBuilderCustomizer, PulsarAdminBuilder> customizers = Customizers
|
||||
.of(PulsarAdminBuilder.class, PulsarAdminBuilderCustomizer::customize);
|
||||
assertThat(customizers.fromField(pulsarAdmin, "adminCustomizers"))
|
||||
.callsInOrder(PulsarAdminBuilder::serviceHttpUrl, "fromPropsCustomizer");
|
||||
});
|
||||
}
|
||||
|
||||
@Test
|
||||
void whenHasUserDefinedCustomizersAppliesInCorrectOrder() {
|
||||
PulsarConnectionDetails connectionDetails = mock(PulsarConnectionDetails.class);
|
||||
given(connectionDetails.getPulsarAdminUrl()).willReturn("fromConnectionDetailsCustomizer");
|
||||
this.contextRunner.withUserConfiguration(PulsarAdminBuilderCustomizersConfig.class)
|
||||
.withBean(PulsarConnectionDetails.class, () -> connectionDetails)
|
||||
.withPropertyValues("spring.pulsar.admin.service-url=fromPropsCustomizer")
|
||||
.run((context) -> {
|
||||
PulsarAdministration pulsarAdmin = context.getBean(PulsarAdministration.class);
|
||||
Customizers<PulsarAdminBuilderCustomizer, PulsarAdminBuilder> customizers = Customizers
|
||||
.of(PulsarAdminBuilder.class, PulsarAdminBuilderCustomizer::customize);
|
||||
assertThat(customizers.fromField(pulsarAdmin, "adminCustomizers")).callsInOrder(
|
||||
PulsarAdminBuilder::serviceHttpUrl, "fromPropsCustomizer", "fromCustomizer1",
|
||||
"fromCustomizer2");
|
||||
PulsarAdminBuilder::serviceHttpUrl, "fromPropsCustomizer",
|
||||
"fromConnectionDetailsCustomizer", "fromCustomizer1", "fromCustomizer2");
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user