Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Properties;
import java.util.Set;
Expand All @@ -33,6 +34,7 @@
import lombok.Setter;
import org.apache.bookkeeper.client.api.DigestType;
import org.apache.pulsar.broker.authorization.PulsarAuthorizationProvider;
import org.apache.pulsar.broker.validator.MultipleListenerValidator;
import org.apache.pulsar.common.configuration.Category;
import org.apache.pulsar.common.configuration.FieldContext;
import org.apache.pulsar.common.configuration.PulsarConfiguration;
Expand All @@ -43,6 +45,7 @@
import org.apache.pulsar.common.policies.data.TopicType;
import org.apache.pulsar.common.protocol.Commands;
import org.apache.pulsar.common.sasl.SaslConstants;
import org.apache.pulsar.policies.data.loadbalancer.AdvertisedListener;

/**
* Pulsar service configuration object.
Expand Down Expand Up @@ -2216,4 +2219,5 @@ public int getBrokerDeleteInactiveTopicsMaxInactiveDurationSeconds() {
return brokerDeleteInactiveTopicsMaxInactiveDurationSeconds;
}
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,14 @@
*/
package org.apache.pulsar.broker;

import org.apache.pulsar.broker.validator.MultipleListenerValidator;
import org.apache.pulsar.policies.data.loadbalancer.AdvertisedListener;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.net.InetAddress;
import java.net.UnknownHostException;
import java.util.Map;

import static org.apache.commons.lang3.StringUtils.isBlank;

Expand All @@ -47,4 +50,30 @@ public static String unsafeLocalhostResolve() {
}
}

/**
* Get the address of Broker, first try to get it from AdvertisedAddress.
* If it is not set, try to get the address set by advertisedListener.
* If it is still not set, get it through InetAddress.getLocalHost().
* @return
*/
public static String getAppliedAdvertisedAddress(ServiceConfiguration configuration) {
Map<String, AdvertisedListener> result = MultipleListenerValidator
.validateAndAnalysisAdvertisedListener(configuration);

String advertisedAddress = configuration.getAdvertisedAddress();
if (advertisedAddress != null) {
return advertisedAddress;
}

AdvertisedListener advertisedListener = result.get(configuration.getInternalListenerName());
if (advertisedListener != null) {
String address = advertisedListener.getBrokerServiceUrl().getHost();
if (address != null) {
return address;
}
}

return getDefaultOrConfiguredAddress(advertisedAddress);
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -19,10 +19,13 @@
package org.apache.pulsar.broker.validator;

import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.ServiceConfigurationUtils;
import org.testng.annotations.Test;

import java.util.Optional;

import static org.testng.AssertJUnit.assertEquals;

/**
* testcase for MultipleListenerValidator.
*/
Expand All @@ -37,6 +40,25 @@ public void testAppearTogether() {
MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config);
}

@Test
public void testGetAppliedAdvertised() throws Exception {
ServiceConfiguration config = new ServiceConfiguration();
config.setBrokerServicePortTls(Optional.of(6651));
config.setAdvertisedListeners("internal:pulsar://192.0.0.1:6660, internal:pulsar+ssl://192.0.0.1:6651");
config.setInternalListenerName("internal");
assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), "192.0.0.1");

config = new ServiceConfiguration();
config.setBrokerServicePortTls(Optional.of(6651));
config.setAdvertisedAddress("192.0.0.2");
assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config), "192.0.0.2");

config.setAdvertisedAddress(null);
assertEquals(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config),
ServiceConfigurationUtils.getDefaultOrConfiguredAddress(null));
}


@Test(expectedExceptions = IllegalArgumentException.class)
public void testListenerDuplicate_1() {
ServiceConfiguration config = new ServiceConfiguration();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -271,22 +271,11 @@ public PulsarService(ServiceConfiguration config,
// Validate correctness of configuration
PulsarConfigurationLoader.isComplete(config);
// validate `advertisedAddress`, `advertisedListeners`, `internalListenerName`
Map<String, AdvertisedListener> result =
MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config);
if (result != null) {
this.advertisedListeners = Collections.unmodifiableMap(result);
} else {
this.advertisedListeners = Collections.unmodifiableMap(Collections.emptyMap());
Comment thread
315157973 marked this conversation as resolved.
}
this.advertisedListeners = MultipleListenerValidator.validateAndAnalysisAdvertisedListener(config);
this.advertisedAddress = ServiceConfigurationUtils.getAppliedAdvertisedAddress(config);
state = State.Init;
// use `internalListenerName` listener as `advertisedAddress`
this.bindAddress = ServiceConfigurationUtils.getDefaultOrConfiguredAddress(config.getBindAddress());
if (!this.advertisedListeners.isEmpty()) {
this.advertisedAddress = this.advertisedListeners.get(
config.getInternalListenerName()).getBrokerServiceUrl().getHost();
} else {
this.advertisedAddress = advertisedAddress(config);
}
this.brokerVersion = PulsarVersion.getVersion();
this.config = config;
this.shutdownService = new MessagingServiceShutdownHook(this, processTerminator);
Expand Down Expand Up @@ -1326,18 +1315,10 @@ public ShutdownService getShutdownService() {
return shutdownService;
}

/**
* Advertised service address.
*
* @return Hostname or IP address the service advertises to the outside world.
*/
public static String advertisedAddress(ServiceConfiguration config) {
return ServiceConfigurationUtils.getDefaultOrConfiguredAddress(config.getAdvertisedAddress());
}

private String brokerUrl(ServiceConfiguration config) {
protected String brokerUrl(ServiceConfiguration config) {
if (config.getBrokerServicePort().isPresent()) {
return brokerUrl(advertisedAddress(config), getBrokerListenPort().get());
return brokerUrl(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config),
getBrokerListenPort().get());
} else {
return null;
}
Expand All @@ -1349,7 +1330,8 @@ public static String brokerUrl(String host, int port) {

public String brokerUrlTls(ServiceConfiguration config) {
if (config.getBrokerServicePortTls().isPresent()) {
return brokerUrlTls(advertisedAddress(config), getBrokerListenPortTls().get());
return brokerUrlTls(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config),
getBrokerListenPortTls().get());
} else {
return null;
}
Expand All @@ -1361,7 +1343,8 @@ public static String brokerUrlTls(String host, int port) {

public String webAddress(ServiceConfiguration config) {
if (config.getWebServicePort().isPresent()) {
return webAddress(advertisedAddress(config), getListenPortHTTP().get());
return webAddress(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config),
getListenPortHTTP().get());
} else {
return null;
}
Expand All @@ -1373,7 +1356,8 @@ public static String webAddress(String host, int port) {

public String webAddressTls(ServiceConfiguration config) {
if (config.getWebServicePortTls().isPresent()) {
return webAddressTls(advertisedAddress(config), getListenPortHTTPS().get());
return webAddressTls(ServiceConfigurationUtils.getAppliedAdvertisedAddress(config),
getListenPortHTTPS().get());
} else {
return null;
}
Expand Down Expand Up @@ -1514,7 +1498,7 @@ public static WorkerConfig initializeWorkerConfigFromBrokerConfig(ServiceConfigu

// worker talks to local broker
String hostname = ServiceConfigurationUtils.getDefaultOrConfiguredAddress(
brokerConfig.getAdvertisedAddress());
ServiceConfigurationUtils.getAppliedAdvertisedAddress(brokerConfig));
workerConfig.setWorkerHostname(hostname);
// inherit broker authorization setting
workerConfig.setAuthenticationEnabled(brokerConfig.isAuthenticationEnabled());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@
import org.apache.pulsar.broker.PulsarServerException;
import org.apache.pulsar.broker.PulsarService;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.ServiceConfigurationUtils;
import org.apache.pulsar.broker.TimeAverageBrokerData;
import org.apache.pulsar.broker.TimeAverageMessageData;
import org.apache.pulsar.broker.loadbalance.BrokerFilter;
Expand Down Expand Up @@ -974,7 +975,7 @@ private void updateLoadBalancingMetrics(final SystemResourceUsage systemResource
List<Metrics> metrics = Lists.newArrayList();
Map<String, String> dimensions = new HashMap<>();

dimensions.put("broker", conf.getAdvertisedAddress());
dimensions.put("broker", ServiceConfigurationUtils.getAppliedAdvertisedAddress(conf));
dimensions.put("metric", "loadBalancing");

Metrics m = Metrics.create(dimensions);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
import org.apache.pulsar.broker.BookKeeperClientFactoryImpl;
import org.apache.pulsar.broker.PulsarService;
import org.apache.pulsar.broker.ServiceConfiguration;
import org.apache.pulsar.broker.ServiceConfigurationUtils;
import org.apache.pulsar.client.api.ClientBuilder;
import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.common.configuration.PulsarConfigurationLoader;
Expand Down Expand Up @@ -100,13 +101,15 @@ public static void main(String[] args) throws Exception {
log.info("Found `brokerServicePortTls` in configuration file. \n"
+ "Will connect pulsar use TLS.");
clientBuilder
.serviceUrl(PulsarService.brokerUrlTls(PulsarService.advertisedAddress(brokerConfig),
.serviceUrl(PulsarService.brokerUrlTls(ServiceConfigurationUtils
.getAppliedAdvertisedAddress(brokerConfig),
brokerConfig.getBrokerServicePortTls().get()))
.allowTlsInsecureConnection(brokerConfig.isTlsAllowInsecureConnection())
.tlsTrustCertsFilePath(brokerConfig.getTlsCertificateFilePath());

} else {
clientBuilder.serviceUrl(PulsarService.brokerUrl(PulsarService.advertisedAddress(brokerConfig),
clientBuilder.serviceUrl(PulsarService.brokerUrl(ServiceConfigurationUtils
.getAppliedAdvertisedAddress(brokerConfig),
brokerConfig.getBrokerServicePort().get()));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,17 +21,48 @@
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.spy;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertNull;
import static org.testng.AssertJUnit.assertSame;
import java.util.Optional;
import lombok.Cleanup;
import lombok.extern.slf4j.Slf4j;
import org.apache.pulsar.broker.auth.MockedPulsarServiceBaseTest;
import org.apache.pulsar.functions.worker.WorkerConfig;
import org.apache.pulsar.functions.worker.WorkerService;
import org.testng.AssertJUnit;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.Test;


@Slf4j
public class PulsarServiceTest {
public class PulsarServiceTest extends MockedPulsarServiceBaseTest {

private boolean useListenerName = false;

@Override
protected void setup() throws Exception {
super.internalSetup();
}

@AfterMethod(alwaysRun = true)
@Override
protected void cleanup() throws Exception {
super.internalCleanup();
useListenerName = false;
resetConfig();
}

@Override
protected void doInitConf() throws Exception {
super.doInitConf();
if (useListenerName) {
conf.setAdvertisedAddress(null);
conf.setBrokerServicePortTls(Optional.of(6651));
conf.setBrokerServicePort(Optional.of(6660));
conf.setWebServicePort(Optional.of(8081));
conf.setWebServicePortTls(Optional.of(8082));
}
}

@Test
public void testGetWorkerService() throws Exception {
Expand Down Expand Up @@ -71,4 +102,21 @@ public void testGetWorkerServiceException() throws Exception {
assertEquals(e.getMessage(), errorMessage);
}
}

@Test
public void testAppliedAdvertised() throws Exception {
useListenerName = true;
conf.setAdvertisedListeners("internal:pulsar://127.0.0.1, internal:pulsar+ssl://127.0.0.1");
conf.setInternalListenerName("internal");
setup();

AssertJUnit.assertEquals(pulsar.getAdvertisedAddress(), "127.0.0.1");
assertNull(pulsar.getConfiguration().getAdvertisedAddress());
assertEquals(conf, pulsar.getConfiguration());
assertEquals(pulsar.brokerUrlTls(conf), "pulsar+ssl://127.0.0.1:6651");
assertEquals(pulsar.brokerUrl(conf), "pulsar://127.0.0.1:6660");
assertEquals(pulsar.webAddress(conf), "http://127.0.0.1:8081");
assertEquals(pulsar.webAddressTls(conf), "https://127.0.0.1:8082");
}

}