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 @@ -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
Comment thread
Technoboy- marked this conversation as resolved.
// 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
Expand Down Expand Up @@ -2746,13 +2745,23 @@ public void cancelPublishBufferLimiting() {
}

private <T> ServerError getErrorCode(CompletableFuture<T> future) {
return getErrorCodeWithErrorLog(future, false, null);
}

private <T> ServerError getErrorCodeWithErrorLog(CompletableFuture<T> future, boolean logIfError,
String errorMessageIfLog) {
ServerError error = ServerError.UnknownError;
try {
future.getNow(null);
} catch (Exception e) {
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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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<Long> ignoreArgumentCaptor = ArgumentCaptor.forClass(Long.class);
final ArgumentCaptor<CompletableFuture> 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<Consumer>(){

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<ConcurrentLongHashMap> 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<Object> 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<ConcurrentLongHashMap> 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<Object> 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();
}
}
}