diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index cbafce14b0c82..ac4a11c1fd05f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -1001,32 +1001,31 @@ protected void handleSubscribe(final CommandSubscribe subscribe) { } if (existingConsumerFuture != null) { - if (existingConsumerFuture.isDone() && !existingConsumerFuture.isCompletedExceptionally()) { - Consumer consumer = existingConsumerFuture.getNow(null); - log.info("[{}] Consumer with the same id is already created:" - + " consumerId={}, consumer={}", - remoteAddress, consumerId, consumer); - commandSender.sendSuccessResponse(requestId); - return null; - } else { + if (!existingConsumerFuture.isDone()){ // There was an early request to create a consumer with same consumerId. This can happen // when // client timeout is lower the broker timeouts. We need to wait until the previous // consumer // creation request either complete or fails. log.warn("[{}][{}][{}] Consumer with id is already present on the connection," - + " consumerId={}", remoteAddress, topicName, subscriptionName, consumerId); - ServerError error = null; - if (!existingConsumerFuture.isDone()) { - error = ServerError.ServiceNotReady; - } else { - error = getErrorCode(existingConsumerFuture); - consumers.remove(consumerId, existingConsumerFuture); - } - commandSender.sendErrorResponse(requestId, error, + + " consumerId={}", remoteAddress, topicName, subscriptionName, consumerId); + commandSender.sendErrorResponse(requestId, ServerError.ServiceNotReady, "Consumer is already present on the connection"); - return null; + } else if (existingConsumerFuture.isCompletedExceptionally()){ + ServerError error = getErrorCodeWithErrorLog(existingConsumerFuture, true, + String.format("Consumer subscribe failure. remoteAddress: %s, subscription: %s", + remoteAddress, subscriptionName)); + consumers.remove(consumerId, existingConsumerFuture); + commandSender.sendErrorResponse(requestId, error, + "Consumer that failed is already present on the connection"); + } else { + Consumer consumer = existingConsumerFuture.getNow(null); + log.info("[{}] Consumer with the same id is already created:" + + " consumerId={}, consumer={}", + remoteAddress, consumerId, consumer); + commandSender.sendSuccessResponse(requestId); } + return null; } boolean createTopicIfDoesNotExist = forceTopicCreation @@ -2746,6 +2745,11 @@ public void cancelPublishBufferLimiting() { } private ServerError getErrorCode(CompletableFuture future) { + return getErrorCodeWithErrorLog(future, false, null); + } + + private ServerError getErrorCodeWithErrorLog(CompletableFuture future, boolean logIfError, + String errorMessageIfLog) { ServerError error = ServerError.UnknownError; try { future.getNow(null); @@ -2753,6 +2757,11 @@ private ServerError getErrorCode(CompletableFuture future) { if (e.getCause() instanceof BrokerServiceException) { error = BrokerServiceException.getClientErrorCode(e.getCause()); } + if (logIfError){ + String finalErrorMessage = StringUtils.isNotBlank(errorMessageIfLog) + ? errorMessageIfLog : "Unknown Error"; + log.error(finalErrorMessage, e); + } } return error; } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java index 4e6db2ee0f3c1..e8fa8c6063923 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ServerCnxTest.java @@ -53,6 +53,8 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Predicate; import java.util.function.Supplier; import org.apache.bookkeeper.common.util.OrderedExecutor; import org.apache.bookkeeper.mledger.AsyncCallbacks.AddEntryCallback; @@ -113,10 +115,13 @@ import org.apache.pulsar.common.protocol.PulsarHandler; import org.apache.pulsar.common.topics.TopicList; import org.apache.pulsar.common.util.FutureUtil; +import org.apache.pulsar.common.util.collections.ConcurrentLongHashMap; import org.apache.pulsar.metadata.api.extended.MetadataStoreExtended; import org.apache.pulsar.metadata.impl.ZKMetadataStore; import org.apache.zookeeper.ZooKeeper; import org.awaitility.Awaitility; +import org.mockito.ArgumentCaptor; +import org.mockito.MockedStatic; import org.mockito.Mockito; import org.mockito.invocation.InvocationOnMock; import org.mockito.stubbing.Answer; @@ -1954,4 +1959,146 @@ public void testGetTopicsOfNamespaceNoChange() throws Exception { channel.finish(); } + + @Test + public void testNeverDelayConsumerFutureWhenNotFail() throws Exception{ + // Mock ServerCnx.field: consumers + ConcurrentLongHashMap.Builder mapBuilder = Mockito.mock(ConcurrentLongHashMap.Builder.class); + Mockito.when(mapBuilder.expectedItems(Mockito.anyInt())).thenReturn(mapBuilder); + Mockito.when(mapBuilder.concurrencyLevel(Mockito.anyInt())).thenReturn(mapBuilder); + ConcurrentLongHashMap consumers = Mockito.mock(ConcurrentLongHashMap.class); + Mockito.when(mapBuilder.build()).thenReturn(consumers); + ArgumentCaptor ignoreArgumentCaptor = ArgumentCaptor.forClass(Long.class); + final ArgumentCaptor deleteTimesMark = ArgumentCaptor.forClass(CompletableFuture.class); + Mockito.when(consumers.remove(ignoreArgumentCaptor.capture())).thenReturn(true); + Mockito.when(consumers.remove(ignoreArgumentCaptor.capture(), deleteTimesMark.capture())).thenReturn(true); + // case1: exists existingConsumerFuture, already complete or delay done after execute 'isDone()' many times + // case2: exists existingConsumerFuture, delay complete after execute 'isDone()' many times + // Why is the design so complicated, see: https://github.com/apache/pulsar/pull/15051 + // Try a delay of 3 stages. The simulation is successful after repeated judgments. + for(AtomicInteger futureWillDoneAfterDelayTimes = new AtomicInteger(1); + futureWillDoneAfterDelayTimes.intValue() <= 3; + futureWillDoneAfterDelayTimes.incrementAndGet()){ + final AtomicInteger futureCallTimes = new AtomicInteger(); + final Consumer mockConsumer = Mockito.mock(Consumer.class); + CompletableFuture existingConsumerFuture = new CompletableFuture(){ + + private boolean complete; + + // delay complete after execute 'isDone()' many times + @Override + public boolean isDone() { + if (complete) { + return true; + } + int executeIsDoneCommandTimes = futureCallTimes.incrementAndGet(); + return executeIsDoneCommandTimes >= futureWillDoneAfterDelayTimes.intValue(); + } + + // if trig "getNow()", then complete + @Override + public Consumer get(){ + complete = true; + return mockConsumer; + } + + // if trig "get()", then complete + @Override + public Consumer get(long timeout, TimeUnit unit){ + complete = true; + return mockConsumer; + } + + // if trig "get()", then complete + @Override + public Consumer getNow(Consumer ifAbsent){ + complete = true; + return mockConsumer; + } + + // never fail + public boolean isCompletedExceptionally(){ + return false; + } + }; + Mockito.when(consumers.putIfAbsent(Mockito.anyLong(), Mockito.any())).thenReturn(existingConsumerFuture); + // do test: delay complete after execute 'isDone()' many times + // Why is the design so complicated, see: https://github.com/apache/pulsar/pull/15051 + try (MockedStatic theMock = Mockito.mockStatic(ConcurrentLongHashMap.class)) { + // Inject consumers to ServerCnx + theMock.when(ConcurrentLongHashMap::newBuilder).thenReturn(mapBuilder); + // reset channels( serverChannel, clientChannel ) + resetChannel(); + setChannelConnected(); + // auth check disable + doReturn(false).when(brokerService).isAuthenticationEnabled(); + doReturn(false).when(brokerService).isAuthorizationEnabled(); + // do subscribe + ByteBuf clientCommand = Commands.newSubscribe(successTopicName, // + successSubName, 1 /* consumer id */, 1 /* request id */, SubType.Exclusive, 0, + "test" /* consumer name */, 0 /* avoid reseting cursor */); + channel.writeInbound(clientCommand); + Object responseObj = getResponse(); + Predicate responseAssert = obj -> { + if (responseObj instanceof CommandSuccess) { + return true; + } + if (responseObj instanceof CommandError) { + CommandError commandError = (CommandError) responseObj; + return ServerError.ServiceNotReady == commandError.getError(); + } + return false; + }; + // assert no consumer-delete event occur + assertFalse(deleteTimesMark.getAllValues().contains(existingConsumerFuture)); + // assert without another error occur + assertTrue(responseAssert.test(responseAssert)); + // Server will not close the connection + assertTrue(channel.isOpen()); + channel.finish(); + } + } + // case3: exists existingConsumerFuture, already complete and exception + CompletableFuture existingConsumerFuture = Mockito.mock(CompletableFuture.class); + Mockito.when(consumers.putIfAbsent(Mockito.anyLong(), Mockito.any())).thenReturn(existingConsumerFuture); + // make consumerFuture delay finish + Mockito.when(existingConsumerFuture.isDone()).thenReturn(true); + // when sync get return, future will return success value. + Mockito.when(existingConsumerFuture.get()).thenThrow(new NullPointerException()); + Mockito.when(existingConsumerFuture.get(Mockito.anyLong(), Mockito.any())). + thenThrow(new NullPointerException()); + Mockito.when(existingConsumerFuture.isCompletedExceptionally()).thenReturn(true); + Mockito.when(existingConsumerFuture.getNow(Mockito.any())).thenThrow(new NullPointerException()); + try (MockedStatic theMock = Mockito.mockStatic(ConcurrentLongHashMap.class)) { + // Inject consumers to ServerCnx + theMock.when(ConcurrentLongHashMap::newBuilder).thenReturn(mapBuilder); + // reset channels( serverChannel, clientChannel ) + resetChannel(); + setChannelConnected(); + // auth check disable + doReturn(false).when(brokerService).isAuthenticationEnabled(); + doReturn(false).when(brokerService).isAuthorizationEnabled(); + // do subscribe + ByteBuf clientCommand = Commands.newSubscribe(successTopicName, // + successSubName, 1 /* consumer id */, 1 /* request id */, SubType.Exclusive, 0, + "test" /* consumer name */, 0 /* avoid reseting cursor */); + channel.writeInbound(clientCommand); + Object responseObj = getResponse(); + Predicate responseAssert = obj -> { + if (responseObj instanceof CommandError) { + CommandError commandError = (CommandError) responseObj; + return ServerError.ServiceNotReady != commandError.getError(); + } + return false; + }; + // assert error response + assertTrue(responseAssert.test(responseAssert)); + // assert consumer-delete event occur + assertEquals(1L, + deleteTimesMark.getAllValues().stream().filter(f -> f == existingConsumerFuture).count()); + // Server will not close the connection + assertTrue(channel.isOpen()); + channel.finish(); + } + } }