From f3a7a69f5fffad6edabe26b4dc13909ca0274e37 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 30 Nov 2021 11:51:19 +0100 Subject: [PATCH 1/5] Transactions: implement support for JWT token authentication between brokers --- .../handlers/kop/KafkaCommandDecoder.java | 6 +- .../handlers/kop/KafkaProtocolHandler.java | 8 +- .../transaction/TransactionCoordinator.java | 10 +- .../TransactionMarkerChannelHandler.java | 176 +++++++++++++++++- .../TransactionMarkerChannelInitializer.java | 2 - .../TransactionMarkerChannelManager.java | 35 +++- .../kop/KopProtocolHandlerTestBase.java | 27 ++- .../handlers/kop/SaslPlainTestBase.java | 110 +++++++++++ 8 files changed, 353 insertions(+), 21 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java index dfefee8436..5aa544ae1b 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaCommandDecoder.java @@ -567,7 +567,7 @@ protected abstract void channelPrepare(ChannelHandlerContext ctx, protected abstract void handleCreatePartitions(KafkaHeaderAndRequest kafkaHeaderAndRequest, CompletableFuture response); - static class KafkaHeaderAndRequest implements Closeable { + public static class KafkaHeaderAndRequest implements Closeable { private static final String DEFAULT_CLIENT_HOST = ""; @@ -576,7 +576,7 @@ static class KafkaHeaderAndRequest implements Closeable { private final ByteBuf buffer; private final SocketAddress remoteAddress; - KafkaHeaderAndRequest(RequestHeader header, + public KafkaHeaderAndRequest(RequestHeader header, AbstractRequest request, ByteBuf buffer, SocketAddress remoteAddress) { @@ -630,7 +630,7 @@ private static boolean isUnsupportedApiVersionsRequest(RequestHeader header) { return header.apiKey() == API_VERSIONS && !API_VERSIONS.isVersionSupported(header.apiVersion()); } - static class KafkaHeaderAndResponse implements Closeable { + public static class KafkaHeaderAndResponse implements Closeable { private final short apiVersion; private final ResponseHeader header; private final AbstractResponse response; diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index ea50e99971..aecf1d9135 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -719,11 +719,17 @@ public TransactionCoordinator initTransactionCoordinator(String tenant, PulsarAd MetadataUtils.createTxnMetadataIfMissing(tenant, pulsarAdmin, clusterData, kafkaConfig); TransactionCoordinator transactionCoordinator = TransactionCoordinator.of( + tenant, + kafkaConfig, transactionConfig, txnTopicClient, brokerService.getPulsar().getLocalMetadataStore(), kopBrokerLookupManager, - OrderedScheduler.newSchedulerBuilder().name("transaction-log-manager").numThreads(1).build(), + OrderedScheduler + .newSchedulerBuilder() + .name("transaction-log-manager-"+tenant) + .numThreads(1) + .build(), Time.SYSTEM, namespacePrefixForMetadata, namespacePrefixForUserTopics); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java index f5f5262571..18f3237d78 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java @@ -22,6 +22,7 @@ import com.google.common.annotations.VisibleForTesting; import com.google.common.collect.Sets; +import io.streamnative.pulsar.handlers.kop.KafkaServiceConfiguration; import io.streamnative.pulsar.handlers.kop.KopBrokerLookupManager; import io.streamnative.pulsar.handlers.kop.SystemTopicClient; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionMetadata.TxnTransitMetadata; @@ -130,19 +131,21 @@ protected TransactionCoordinator(TransactionConfig transactionConfig, this.time = time; } - public static TransactionCoordinator of(TransactionConfig transactionConfig, + public static TransactionCoordinator of(String tenant, + KafkaServiceConfiguration kafkaConfig, + TransactionConfig transactionConfig, SystemTopicClient txnTopicClient, MetadataStoreExtended metadataStore, KopBrokerLookupManager kopBrokerLookupManager, ScheduledExecutorService scheduler, Time time, String namespacePrefixForMetadata, - String namespacePrefixForUserTopics) { + String namespacePrefixForUserTopics) throws Exception { TransactionStateManager transactionStateManager = new TransactionStateManager(transactionConfig, txnTopicClient, scheduler, time); return new TransactionCoordinator( transactionConfig, - new TransactionMarkerChannelManager(null, transactionStateManager, + new TransactionMarkerChannelManager(tenant, kafkaConfig, transactionStateManager, kopBrokerLookupManager, false, namespacePrefixForUserTopics), scheduler, new ProducerIdManager(transactionConfig.getBrokerId(), metadataStore), @@ -991,6 +994,7 @@ public void shutdown() { producerIdManager.shutdown(); txnManager.shutdown(); transactionMarkerChannelManager.close(); + scheduler.shutdown(); // TODO shutdown txn log.info("Shutdown transaction coordinator complete."); } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java index bb0aaa35e0..d99919c697 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java @@ -13,28 +13,40 @@ */ package io.streamnative.pulsar.handlers.kop.coordinator.transaction; +import static java.nio.charset.StandardCharsets.UTF_8; import static org.apache.kafka.common.protocol.Errors.REQUEST_TIMED_OUT; import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; +import io.streamnative.pulsar.handlers.kop.KafkaCommandDecoder; import java.net.InetSocketAddress; import java.nio.ByteBuffer; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; import java.util.concurrent.atomic.AtomicInteger; +import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.util.collections.ConcurrentLongHashMap; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.protocol.ApiKeys; import org.apache.kafka.common.protocol.Errors; +import org.apache.kafka.common.protocol.types.Struct; +import org.apache.kafka.common.requests.AbstractResponse; import org.apache.kafka.common.requests.KopRequestUtils; import org.apache.kafka.common.requests.RequestHeader; import org.apache.kafka.common.requests.ResponseHeader; +import org.apache.kafka.common.requests.SaslAuthenticateRequest; +import org.apache.kafka.common.requests.SaslAuthenticateResponse; +import org.apache.kafka.common.requests.SaslHandshakeRequest; import org.apache.kafka.common.requests.WriteTxnMarkersRequest; import org.apache.kafka.common.requests.WriteTxnMarkersResponse; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.common.util.FutureUtil; /** @@ -45,6 +57,8 @@ public class TransactionMarkerChannelHandler extends ChannelInboundHandlerAdapte private final CompletableFuture cnx = new CompletableFuture<>(); private final ConcurrentLongHashMap inFlightRequestMap = new ConcurrentLongHashMap<>(); + private final ConcurrentLongHashMap genericRequestMap + = new ConcurrentLongHashMap<>(); private final AtomicInteger correlationId = new AtomicInteger(0); private final TransactionMarkerChannelManager transactionMarkerChannelManager; @@ -68,6 +82,13 @@ public void enqueueRequest(WriteTxnMarkersRequest request, }); } + @AllArgsConstructor + private static final class PendingGenericRequest { + CompletableFuture response; + ApiKeys apiKeys; + short apiVersion; + } + private class InFlightRequest { private final long requestId; @@ -122,18 +143,23 @@ public void channelActive(ChannelHandlerContext channelHandlerContext) throws Ex log.debug("channelActive"); } log.info("[TransactionMarkerChannelHandler] channelActive to {}", channelHandlerContext.channel()); - this.cnx.complete(channelHandlerContext); + handleAuthentication(channelHandlerContext); super.channelActive(channelHandlerContext); } @Override public void channelInactive(ChannelHandlerContext channelHandlerContext) throws Exception { - log.info("[TransactionMarkerChannelHandler] channelInactive, failing {} pending requests", - inFlightRequestMap.size()); + log.info("[TransactionMarkerChannelHandler] channelInactive, failing {} + {} pending requests", + inFlightRequestMap.size(), genericRequestMap.size()); + final Exception exception = new Exception("Connection to remote broker closed"); inFlightRequestMap.forEach((k, v) -> { - v.onError(new Exception("Connection to remote broker closed")); + v.onError(exception); }); inFlightRequestMap.clear(); + genericRequestMap.forEach((k,v)-> { + v.response.completeExceptionally(exception); + }); + genericRequestMap.clear(); transactionMarkerChannelManager.channelFailed((InetSocketAddress) channelHandlerContext .channel() .remoteAddress(), this); @@ -145,20 +171,33 @@ public void channelRead(ChannelHandlerContext channelHandlerContext, Object o) t ByteBuffer nio = ((ByteBuf) o).nioBuffer(); ResponseHeader responseHeader = ResponseHeader.parse(nio); InFlightRequest inFlightRequest = inFlightRequestMap.remove(responseHeader.correlationId()); - if (inFlightRequest == null) { - log.error("Miss the inFlightRequest with correlationId {}.", responseHeader.correlationId()); + if (inFlightRequest != null) { + inFlightRequest.onComplete(nio); + return; + } + PendingGenericRequest genericRequest = genericRequestMap.remove(responseHeader.correlationId()); + if (genericRequest != null) { + Struct responseBody = genericRequest.apiKeys.parseResponse(genericRequest.apiVersion, nio); + AbstractResponse response = AbstractResponse.parseResponse(genericRequest.apiKeys, responseBody); + genericRequest.response.complete(response); return; } - inFlightRequest.onComplete(nio); + log.error("Miss the inFlightRequest with correlationId {}.", responseHeader.correlationId()); } @Override public void exceptionCaught(ChannelHandlerContext channelHandlerContext, Throwable throwable) throws Exception { log.error("Transaction marker channel handler caught exception.", throwable); + final Exception exception = + new Exception("Transaction marker channel handler caught exception: " + throwable, throwable); inFlightRequestMap.forEach((k, v) -> { - v.onError(new Exception("Transaction marker channel handler caught exception: " + throwable, throwable)); + v.onError(exception); }); inFlightRequestMap.clear(); + genericRequestMap.forEach((k,v)-> { + v.response.completeExceptionally(exception); + }); + genericRequestMap.clear(); channelHandlerContext.close(); } @@ -171,4 +210,125 @@ public void close() { }); } + public void handleAuthentication(ChannelHandlerContext channelHandlerContext) { + if (!transactionMarkerChannelManager.getKafkaConfig().isAuthenticationEnabled()) { + this.cnx.complete(channelHandlerContext); + return; + } + saslHandshake(channelHandlerContext) + .thenCompose(this::authenticate) + .thenApply(cnx::complete) + .exceptionally(err -> { + cnx.completeExceptionally(err); + return null; + }); + } + + private void sendGenericRequestOnTheWire(ChannelHandlerContext channel, + KafkaCommandDecoder.KafkaHeaderAndRequest request, + CompletableFuture result) { + long correlationId = request.getHeader().correlationId(); + genericRequestMap.put(correlationId, new PendingGenericRequest(result, + request.getHeader().apiKey(), + request.getHeader().apiVersion())); + channel.writeAndFlush(request.getBuffer()) + .addListener(writeFuture -> { + if (!writeFuture.isSuccess()) { + genericRequestMap.remove(correlationId); + // cannot write, so we have to "close()" and trigger failure of every other + // pending request and discard the reference to this connection + channel.close(); + result.completeExceptionally(writeFuture.cause()); + } + }); + } + + private CompletableFuture saslHandshake(ChannelHandlerContext channel) { + KafkaCommandDecoder.KafkaHeaderAndRequest fullRequest = buildSASLRequest(); + CompletableFuture result = new CompletableFuture<>(); + sendGenericRequestOnTheWire(channel, fullRequest, result); + result.exceptionally(error -> { + // ensure that we close the channel + channel.close(); + return null; + }); + return result.thenApply(response -> { + log.debug("SASL Handshake completed with success"); + return channel; + }); + } + + private KafkaCommandDecoder.KafkaHeaderAndRequest buildSASLRequest() { + RequestHeader header = new RequestHeader( + ApiKeys.SASL_HANDSHAKE, + ApiKeys.SASL_HANDSHAKE.latestVersion(), + "tx", //ignored + correlationId.incrementAndGet() + ); + SaslHandshakeRequest request = new SaslHandshakeRequest + .Builder("PLAIN") + .build(); + ByteBuffer buffer = request.serialize(header); + KafkaCommandDecoder.KafkaHeaderAndRequest fullRequest = new KafkaCommandDecoder.KafkaHeaderAndRequest( + header, + request, + Unpooled.wrappedBuffer(buffer), + null + ); + return fullRequest; + } + + private CompletableFuture authenticate(final ChannelHandlerContext channel) { + CompletableFuture internal = authenticateInternal(channel); + // ensure that we close the channel + internal.exceptionally(error -> { + channel.close(); + return null; + }); + return internal; + } + + private CompletableFuture authenticateInternal(ChannelHandlerContext channel) { + RequestHeader header = new RequestHeader( + ApiKeys.SASL_AUTHENTICATE, + ApiKeys.SASL_AUTHENTICATE.latestVersion(), + "tx", // ignored + correlationId.incrementAndGet() + ); + String prefix = "TX"; // the prefix TX means nothing, it is ignored by SaslUtils#parseSaslAuthBytes + String authUsername = transactionMarkerChannelManager.getAuthenticationUsername(); + String authPassword = transactionMarkerChannelManager.getAuthenticationPassword(); + String usernamePassword = prefix + + "\u0000" + authUsername + + "\u0000" + authPassword; + byte[] saslAuthBytes = usernamePassword.getBytes(UTF_8); + SaslAuthenticateRequest request = new SaslAuthenticateRequest + .Builder(ByteBuffer.wrap(saslAuthBytes)) + .build(); + + ByteBuffer buffer = request.serialize(header); + + KafkaCommandDecoder.KafkaHeaderAndRequest fullRequest = new KafkaCommandDecoder.KafkaHeaderAndRequest( + header, + request, + Unpooled.wrappedBuffer(buffer), + null + ); + CompletableFuture result = new CompletableFuture<>(); + sendGenericRequestOnTheWire(channel, fullRequest, result); + return result.thenApply(response -> { + SaslAuthenticateResponse saslResponse = (SaslAuthenticateResponse) response; + if (saslResponse.error() != Errors.NONE) { + log.error("Failed authentication against KOP broker {}{}", saslResponse.error(), + saslResponse.errorMessage()); + close(); + throw new CompletionException(saslResponse.error().exception()); + } else { + log.debug("Success step AUTH to KOP broker {} {} {}", saslResponse.error(), + saslResponse.errorMessage(), saslResponse.saslAuthBytes()); + } + return channel; + }); + } + } diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java index 98ea51e849..9462569dc8 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelInitializer.java @@ -30,7 +30,6 @@ */ public class TransactionMarkerChannelInitializer extends ChannelInitializer { - private final KafkaServiceConfiguration kafkaConfig; private final boolean enableTls; private final SslContextFactory.Server sslContextFactory; private final TransactionMarkerChannelManager transactionMarkerChannelManager; @@ -38,7 +37,6 @@ public class TransactionMarkerChannelInitializer extends ChannelInitializer producer = kProducer.getProducer(); + + producer.initTransactions(); + + int totalTxnCount = 10; + int messageCountPerTxn = 10; + + String lastMessage = ""; + for (int txnIndex = 0; txnIndex < totalTxnCount; txnIndex++) { + producer.beginTransaction(); + + String contentBase; + if (txnIndex % 2 != 0) { + contentBase = "commit msg txnIndex %s messageIndex %s"; + } else { + contentBase = "abort msg txnIndex %s messageIndex %s"; + } + + for (int messageIndex = 0; messageIndex < messageCountPerTxn; messageIndex++) { + String msgContent = String.format(contentBase, txnIndex, messageIndex); + log.info("send txn message {}", msgContent); + lastMessage = msgContent; + producer.send(new ProducerRecord<>(topicName, messageIndex, msgContent)).get(); + } + + if (txnIndex % 2 != 0) { + producer.commitTransaction(); + } else { + producer.abortTransaction(); + } + } + + consumeTxnMessage(topicName, totalTxnCount * messageCountPerTxn, lastMessage, isolation); + } + + + private void consumeTxnMessage(String topicName, + int totalMessageCount, + String lastMessage, + String isolation) throws InterruptedException { + + @Cleanup + KConsumer kConsumer = new KConsumer(topicName , "localhost", getKafkaBrokerPort(), false, + TENANT + "/" + NAMESPACE, "token:" + userToken, "consumeTxnMessage-" + UUID.randomUUID(), + IntegerDeserializer.class.getName(), + StringDeserializer.class.getName(), isolation); + + KafkaConsumer consumer = kConsumer.getConsumer(); + consumer.subscribe(Collections.singleton(topicName)); + + log.info("the last message is: {}", lastMessage); + AtomicInteger receiveCount = new AtomicInteger(0); + while (true) { + ConsumerRecords consumerRecords = + consumer.poll(Duration.of(100, ChronoUnit.MILLIS)); + + boolean readFinish = false; + for (ConsumerRecord record : consumerRecords) { + log.info("Fetch for receive record offset: {}, key: {}, value: {}", + record.offset(), record.key(), record.value()); + receiveCount.incrementAndGet(); + if (lastMessage.equalsIgnoreCase(record.value())) { + log.info("receive the last message"); + readFinish = true; + } + } + + if (readFinish) { + log.info("Fetch for read finish."); + break; + } + } + log.info("Fetch for receive message finish. isolation: {}, receive count: {}", isolation, receiveCount.get()); + + if (isolation.equals("read_committed")) { + Assert.assertEquals(receiveCount.get(), totalMessageCount / 2); + } else { + Assert.assertEquals(receiveCount.get(), totalMessageCount); + } + log.info("Fetch for finish consume messages. isolation: {}", isolation); + } } From e12fdd6a93b7258b6001feee484187c230a9ae24 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 30 Nov 2021 12:13:18 +0100 Subject: [PATCH 2/5] Fix build --- .../handlers/kop/KafkaProtocolHandler.java | 16 ++++------------ .../transaction/TransactionCoordinator.java | 7 ++++--- .../TransactionMarkerChannelHandler.java | 9 +++------ .../TransactionMarkerChannelManager.java | 4 +--- .../pulsar/handlers/kop/SaslPlainTestBase.java | 11 ++++++----- 5 files changed, 18 insertions(+), 29 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java index aecf1d9135..b9edc95764 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/KafkaProtocolHandler.java @@ -514,13 +514,9 @@ private TransactionCoordinator createAndBootTransactionCoordinator(String tenant .brokerServiceUrl(brokerService.getPulsar().getBrokerServiceUrl()) .brokerServiceUrlTls(brokerService.getPulsar().getBrokerServiceUrlTls()) .build(); - - String namespacePrefixForMetadata = MetadataUtils.constructMetadataNamespace(tenant, kafkaConfig); - String namespacePrefixForUserTopics = MetadataUtils.constructUserTopicsNamespace(tenant, kafkaConfig); try { TransactionCoordinator transactionCoordinator = - initTransactionCoordinator(tenant, brokerService.getPulsar().getAdminClient(), clusterData, - namespacePrefixForMetadata, namespacePrefixForUserTopics); + initTransactionCoordinator(tenant, brokerService.getPulsar().getAdminClient(), clusterData); // Listening transaction topic load/unload brokerService.pulsar() .getNamespaceService() @@ -703,9 +699,7 @@ protected GroupCoordinator startGroupCoordinator(String tenant, SystemTopicClien } public TransactionCoordinator initTransactionCoordinator(String tenant, PulsarAdmin pulsarAdmin, - ClusterData clusterData, - String namespacePrefixForMetadata, - String namespacePrefixForUserTopics) throws Exception { + ClusterData clusterData) throws Exception { TransactionConfig transactionConfig = TransactionConfig.builder() .transactionLogNumPartitions(kafkaConfig.getTxnLogTopicNumPartitions()) .transactionMetadataTopicName(MetadataUtils.constructTxnLogTopicBaseName(tenant, kafkaConfig)) @@ -727,12 +721,10 @@ public TransactionCoordinator initTransactionCoordinator(String tenant, PulsarAd kopBrokerLookupManager, OrderedScheduler .newSchedulerBuilder() - .name("transaction-log-manager-"+tenant) + .name("transaction-log-manager-" + tenant) .numThreads(1) .build(), - Time.SYSTEM, - namespacePrefixForMetadata, - namespacePrefixForUserTopics); + Time.SYSTEM); transactionCoordinator.startup(kafkaConfig.isEnableTransactionalIdExpiration()).get(); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java index 18f3237d78..f86128d229 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionCoordinator.java @@ -27,6 +27,7 @@ import io.streamnative.pulsar.handlers.kop.SystemTopicClient; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionMetadata.TxnTransitMetadata; import io.streamnative.pulsar.handlers.kop.coordinator.transaction.TransactionStateManager.CoordinatorEpochAndTxnMetadata; +import io.streamnative.pulsar.handlers.kop.utils.MetadataUtils; import io.streamnative.pulsar.handlers.kop.utils.ProducerIdAndEpoch; import java.util.ArrayList; import java.util.HashMap; @@ -138,9 +139,9 @@ public static TransactionCoordinator of(String tenant, MetadataStoreExtended metadataStore, KopBrokerLookupManager kopBrokerLookupManager, ScheduledExecutorService scheduler, - Time time, - String namespacePrefixForMetadata, - String namespacePrefixForUserTopics) throws Exception { + Time time) throws Exception { + String namespacePrefixForMetadata = MetadataUtils.constructMetadataNamespace(tenant, kafkaConfig); + String namespacePrefixForUserTopics = MetadataUtils.constructUserTopicsNamespace(tenant, kafkaConfig); TransactionStateManager transactionStateManager = new TransactionStateManager(transactionConfig, txnTopicClient, scheduler, time); return new TransactionCoordinator( diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java index d99919c697..87117b76af 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java @@ -45,8 +45,6 @@ import org.apache.kafka.common.requests.SaslHandshakeRequest; import org.apache.kafka.common.requests.WriteTxnMarkersRequest; import org.apache.kafka.common.requests.WriteTxnMarkersResponse; -import org.apache.pulsar.client.api.PulsarClientException; -import org.apache.pulsar.common.util.FutureUtil; /** @@ -57,8 +55,7 @@ public class TransactionMarkerChannelHandler extends ChannelInboundHandlerAdapte private final CompletableFuture cnx = new CompletableFuture<>(); private final ConcurrentLongHashMap inFlightRequestMap = new ConcurrentLongHashMap<>(); - private final ConcurrentLongHashMap genericRequestMap - = new ConcurrentLongHashMap<>(); + private final ConcurrentLongHashMap genericRequestMap = new ConcurrentLongHashMap<>(); private final AtomicInteger correlationId = new AtomicInteger(0); private final TransactionMarkerChannelManager transactionMarkerChannelManager; @@ -156,7 +153,7 @@ public void channelInactive(ChannelHandlerContext channelHandlerContext) throws v.onError(exception); }); inFlightRequestMap.clear(); - genericRequestMap.forEach((k,v)-> { + genericRequestMap.forEach((k, v)-> { v.response.completeExceptionally(exception); }); genericRequestMap.clear(); @@ -194,7 +191,7 @@ public void exceptionCaught(ChannelHandlerContext channelHandlerContext, Throwab v.onError(exception); }); inFlightRequestMap.clear(); - genericRequestMap.forEach((k,v)-> { + genericRequestMap.forEach((k, v)-> { v.response.completeExceptionally(exception); }); genericRequestMap.clear(); diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java index 71781f508a..9c70008d43 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelManager.java @@ -47,8 +47,6 @@ import org.apache.kafka.common.requests.TransactionResult; import org.apache.kafka.common.requests.WriteTxnMarkersRequest; import org.apache.kafka.common.requests.WriteTxnMarkersRequest.TxnMarkerEntry; -import org.apache.pulsar.client.api.Authentication; -import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.impl.AuthenticationUtil; import org.apache.pulsar.client.impl.auth.AuthenticationToken; import org.apache.pulsar.common.util.FutureUtil; @@ -497,6 +495,6 @@ public String getAuthenticationPassword() { if (authenticationToken == null) { return ""; } - return "token:"+ authenticationToken; + return "token:" + authenticationToken; } } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java index 8bf48138ab..c69276002b 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java @@ -33,7 +33,6 @@ import javax.crypto.SecretKey; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; -import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -228,12 +227,14 @@ void clientWithoutAuth() throws Exception { @Test public void transactionsReadCommittedTest() throws Exception { - basicProduceAndConsumeTest(TENANT + "/" + NAMESPACE + "/" + "read-committed-test", "txn-11", "read_committed"); + basicProduceAndConsumeTest(TENANT + "/" + NAMESPACE + "/" + "read-committed-test", "txn-11", + "read_committed"); } @Test(timeOut = 1000 * 10) public void transactionsReadUncommittedTest() throws Exception { - basicProduceAndConsumeTest(TENANT + "/" + NAMESPACE + "/" + "read-uncommitted-test", "txn-12", "read_uncommitted"); + basicProduceAndConsumeTest(TENANT + "/" + NAMESPACE + "/" + "read-uncommitted-test", "txn-12", + "read_uncommitted"); } private void basicProduceAndConsumeTest(String topicName, @@ -320,9 +321,9 @@ private void consumeTxnMessage(String topicName, log.info("Fetch for receive message finish. isolation: {}, receive count: {}", isolation, receiveCount.get()); if (isolation.equals("read_committed")) { - Assert.assertEquals(receiveCount.get(), totalMessageCount / 2); + assertEquals(receiveCount.get(), totalMessageCount / 2); } else { - Assert.assertEquals(receiveCount.get(), totalMessageCount); + assertEquals(receiveCount.get(), totalMessageCount); } log.info("Fetch for finish consume messages. isolation: {}", isolation); } From 996d61224c071f54ffb2fd41f8a3a1c14b12c0ae Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 30 Nov 2021 12:21:33 +0100 Subject: [PATCH 3/5] fix checkstyle --- .../io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java | 1 - 1 file changed, 1 deletion(-) diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java index c69276002b..0922378760 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java @@ -52,7 +52,6 @@ import org.apache.pulsar.client.impl.auth.AuthenticationToken; import org.apache.pulsar.common.policies.data.AuthAction; import org.apache.pulsar.common.policies.data.TenantInfo; -import org.testng.Assert; import org.testng.annotations.BeforeClass; import org.testng.annotations.Test; From ae24ceedbf70d73bd4080a220b37f7748a44de38 Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Tue, 30 Nov 2021 12:28:13 +0100 Subject: [PATCH 4/5] Make Spotbugs happy --- .../transaction/TransactionMarkerChannelHandler.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java index 87117b76af..de360b0a3f 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/coordinator/transaction/TransactionMarkerChannelHandler.java @@ -217,7 +217,7 @@ public void handleAuthentication(ChannelHandlerContext channelHandlerContext) { .thenApply(cnx::complete) .exceptionally(err -> { cnx.completeExceptionally(err); - return null; + return false; }); } From 7c22bd230808c7b84e45f9e025aba0b8c77547fc Mon Sep 17 00:00:00 2001 From: Enrico Olivelli Date: Thu, 2 Dec 2021 13:25:46 +0100 Subject: [PATCH 5/5] Fix tests --- .../io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java index 0922378760..45c37e8595 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/kop/SaslPlainTestBase.java @@ -99,6 +99,7 @@ protected void setup() throws Exception { ((KafkaServiceConfiguration) conf).setKafkaMetadataTenant("internal"); ((KafkaServiceConfiguration) conf).setKafkaMetadataNamespace("__kafka"); + conf.setKafkaTransactionCoordinatorEnabled(true); conf.setClusterName(super.configClusterName); conf.setAuthorizationEnabled(true); conf.setAuthenticationEnabled(true);