diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 64e4b7ff4ed49..a1cdbd997c529 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -2986,6 +2986,10 @@ public void publishTxnMessage(TxnID txnID, ByteBuf headersAndPayload, PublishCon decrementPendingWriteOpsAndCheck(); }) .exceptionally(throwable -> { + throwable = throwable.getCause(); + if (!(throwable instanceof ManagedLedgerException)) { + throwable = new ManagedLedgerException(throwable); + } addFailed((ManagedLedgerException) throwable, publishContext); return null; }); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java index 79f7f35732754..1b0fb6697ec3f 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/transaction/buffer/impl/TopicTransactionBuffer.java @@ -219,7 +219,7 @@ public void addComplete(Position position, ByteBuf entryData, Object ctx) { @Override public void addFailed(ManagedLedgerException exception, Object ctx) { log.error("Failed to append buffer to txn {}", txnId, exception); - completableFuture.completeExceptionally(new PersistenceException(exception)); + completableFuture.completeExceptionally(exception); } }, null); return completableFuture; diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 971f8b24bce8b..bb172b599c730 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -18,15 +18,19 @@ */ package org.apache.pulsar.broker.transaction; +import static java.nio.charset.StandardCharsets.UTF_8; import static org.apache.pulsar.transaction.coordinator.impl.MLTransactionLogImpl.TRANSACTION_LOG_PREFIX; import com.google.common.collect.Sets; +import io.netty.buffer.Unpooled; import java.lang.reflect.Field; import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; +import org.apache.bookkeeper.mledger.ManagedLedgerException; import org.apache.bookkeeper.mledger.ManagedLedgerFactory; import org.apache.bookkeeper.mledger.impl.ManagedLedgerFactoryImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; @@ -295,4 +299,73 @@ public void testTakeSnapshotBeforeBuildTxnProducer() throws Exception { Assert.assertEquals(snapshot1.getMaxReadPositionEntryId(), 1); }); } + + + @Test + public void testAppendBufferWithNotManageLedgerExceptionCanCastToMLE() + throws Exception { + String topic = "persistent://pulsar/system/testReCreateTopic"; + admin.topics().createNonPartitionedTopic(topic); + + PersistentTopic persistentTopic = + (PersistentTopic) pulsarServiceList.get(0).getBrokerService() + .getTopic(topic, false) + .get().get(); + CountDownLatch countDownLatch = new CountDownLatch(1); + Topic.PublishContext publishContext = new Topic.PublishContext() { + + @Override + public String getProducerName() { + return "test"; + } + + public long getSequenceId() { + return 30; + } + /** + * Return the producer name for the original producer. + * + * For messages published locally, this will return the same local producer name, though in case of + * replicated messages, the original producer name will differ + */ + public String getOriginalProducerName() { + return "test"; + } + + public long getOriginalSequenceId() { + return 30; + } + + public long getHighestSequenceId() { + return 30; + } + + public long getOriginalHighestSequenceId() { + return 30; + } + + public long getNumberOfMessages() { + return 30; + } + + @Override + public void completed(Exception e, long ledgerId, long entryId) { + Assert.assertTrue(e.getCause() instanceof ManagedLedgerException.ManagedLedgerAlreadyClosedException); + countDownLatch.countDown(); + } + }; + + //Close topic manageLedger. + persistentTopic.getManagedLedger().close(); + + //Publish to a closed managerLedger to test ManagerLedgerException. + persistentTopic.publishTxnMessage(new TxnID(123L, 321L), + Unpooled.copiedBuffer("message", UTF_8), publishContext); + + //If it times out, it means that the assertTrue in publishContext.completed is failed. + Awaitility.await().until(() -> { + countDownLatch.await(); + return true; + }); + } } \ No newline at end of file