diff --git a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java index d6cd589ce2f4e..3ab5890f426e6 100644 --- a/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java +++ b/pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java @@ -750,6 +750,12 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private Set brokerInterceptors = Sets.newTreeSet(); + @FieldContext( + category = CATEGORY_SERVER, + doc = "Enable or disable the broker interceptor, which is only used for testing for now" + ) + private boolean disableBrokerInterceptors = true; + @FieldContext( doc = "There are two policies when zookeeper session expired happens, \"shutdown\" and \"reconnect\". \n\n" + " If uses \"shutdown\" policy, shutdown the broker when zookeeper session expired happens.\n\n" @@ -757,6 +763,7 @@ public class ServiceConfiguration implements PulsarConfiguration { ) private String zookeeperSessionExpiredPolicy = "shutdown"; + /**** --- Messaging Protocols --- ****/ @FieldContext( diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/ResponseHandlerFilter.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/ResponseHandlerFilter.java index 823209ce8ef32..647590f68a9dc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/ResponseHandlerFilter.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/ResponseHandlerFilter.java @@ -43,10 +43,12 @@ public class ResponseHandlerFilter implements Filter { private final String brokerAddress; private final BrokerInterceptor interceptor; + private final boolean interceptorEnabled; public ResponseHandlerFilter(PulsarService pulsar) { this.brokerAddress = pulsar.getAdvertisedAddress(); this.interceptor = pulsar.getBrokerInterceptor(); + this.interceptorEnabled = !pulsar.getConfig().getBrokerInterceptors().isEmpty(); } @Override @@ -63,8 +65,9 @@ public void doFilter(ServletRequest request, ServletResponse response, FilterCha /* connection is already invalidated */ } } - interceptor.onWebserviceResponse(request, response); - + if (interceptorEnabled) { + interceptor.onWebserviceResponse(request, response); + } } @Override diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java index 5386c2ffb35f2..3fd2d50491f09 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/web/WebService.java @@ -155,8 +155,11 @@ public void addServlet(String path, ServletHolder servletHolder, boolean require }); } - context.addFilter(new FilterHolder(new PreInterceptFilter(pulsar.getBrokerInterceptor())), + if (!pulsar.getConfig().getBrokerInterceptors().isEmpty() || !pulsar.getConfig().isDisableBrokerInterceptors()) { + // Enable PreInterceptFilter only when interceptors are enabled + context.addFilter(new FilterHolder(new PreInterceptFilter(pulsar.getBrokerInterceptor())), MATCH_ALL, EnumSet.allOf(DispatcherType.class)); + } if (requiresAuthentication && pulsar.getConfiguration().isAuthenticationEnabled()) { FilterHolder filter = new FilterHolder(new AuthenticationFilter( diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorTest.java index 55edf494c7a6f..5ff73fc0b0005 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/intercept/BrokerInterceptorTest.java @@ -51,6 +51,8 @@ public class BrokerInterceptorTest extends ProducerConsumerBase { @BeforeMethod public void setup() throws Exception { + this.conf.setDisableBrokerInterceptors(false); + this.listener1 = mock(BrokerInterceptor.class); this.ncl1 = mock(NarClassLoader.class); this.listener2 = mock(BrokerInterceptor.class);