From 471d4aa32504e4f302d5d096f5ff64827eec9c57 Mon Sep 17 00:00:00 2001 From: congbobo184 Date: Fri, 2 Sep 2022 13:33:04 +0800 Subject: [PATCH 1/4] [feat][txn] Add getState in transaction for client API --- .../broker/transaction/TransactionTest.java | 48 +++++++++++++++++++ .../client/api/transaction/Transaction.java | 18 +++++++ .../impl/transaction/TransactionImpl.java | 15 ++---- 3 files changed, 71 insertions(+), 10 deletions(-) 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 d0674721c00be..49af052b46ae8 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 @@ -1469,4 +1469,52 @@ public Object answer(InvocationOnMock invocation) throws Throwable { Assert.assertTrue(t instanceof BrokerServiceException.ServiceUnitNotReadyException); } } + + @Test + public void testGetTxnState() throws Exception { + Transaction transaction = pulsarClient.newTransaction().withTransactionTimeout(1, TimeUnit.SECONDS) + .build().get(); + + // test OPEN and TIMEOUT + assertEquals(transaction.getTxnState(), Transaction.State.OPEN); + Transaction timeoutTxn = transaction; + Awaitility.await().until(() -> timeoutTxn.getTxnState() == Transaction.State.TIMEOUT); + + // test abort + transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) + .build().get(); + transaction.abort().get(); + assertEquals(transaction.getTxnState(), Transaction.State.ABORTED); + + // test commit + transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) + .build().get(); + transaction.commit().get(); + assertEquals(transaction.getTxnState(), Transaction.State.COMMITTED); + + // test error + transaction = pulsarClient.newTransaction().withTransactionTimeout(1, TimeUnit.SECONDS) + .build().get(); + pulsarServiceList.get(0).getTransactionMetadataStoreService() + .endTransaction(transaction.getTxnID(), 0, false); + transaction.commit(); + Transaction errorTxn = transaction; + Awaitility.await().until(() -> errorTxn.getTxnState() == Transaction.State.ERROR); + + // test committing + transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) + .build().get(); + ((TransactionImpl) transaction).registerSendOp(new CompletableFuture<>()); + transaction.commit(); + Transaction committingTxn = transaction; + Awaitility.await().until(() -> committingTxn.getTxnState() == Transaction.State.COMMITTING); + + // test aborting + transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) + .build().get(); + ((TransactionImpl) transaction).registerSendOp(new CompletableFuture<>()); + transaction.abort(); + Transaction abortingTxn = transaction; + Awaitility.await().until(() -> abortingTxn.getTxnState() == Transaction.State.ABORTING); + } } diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java index fd4cf0bc1665c..71cab13181bff 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java @@ -29,6 +29,16 @@ @InterfaceStability.Evolving public interface Transaction { + enum State { + OPEN, + COMMITTING, + ABORTING, + COMMITTED, + ABORTED, + ERROR, + TIMEOUT + } + /** * Commit the transaction. * @@ -48,4 +58,12 @@ public interface Transaction { * @return {@link TxnID} the txnID. */ TxnID getTxnID(); + + /** + * Get transaction state. + * + * @return {@link State} the state of the transaction. + */ + State getTxnState(); + } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java index 55b20438693e3..0c0bcca07d623 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java @@ -73,16 +73,6 @@ public void run(Timeout timeout) throws Exception { STATE_UPDATE.compareAndSet(this, State.OPEN, State.TIMEOUT); } - public enum State { - OPEN, - COMMITTING, - ABORTING, - COMMITTED, - ABORTED, - ERROR, - TIMEOUT - } - TransactionImpl(PulsarClientImpl client, long transactionTimeoutMs, long txnIdLeastBits, @@ -215,6 +205,11 @@ public TxnID getTxnID() { return new TxnID(txnIdMostBits, txnIdLeastBits); } + @Override + public State getTxnState() { + return state; + } + public boolean checkIfOpen(CompletableFuture completableFuture) { if (state == State.OPEN) { return true; From e2b9c161c0d7166c2c08edd9b2dc6109789e93f0 Mon Sep 17 00:00:00 2001 From: congbobo184 Date: Fri, 2 Sep 2022 15:24:02 +0800 Subject: [PATCH 2/4] fix some comments --- .../broker/transaction/TransactionTest.java | 14 +++---- .../client/impl/TransactionEndToEndTest.java | 2 +- .../client/api/transaction/Transaction.java | 41 ++++++++++++++++++- .../impl/transaction/TransactionImpl.java | 4 +- 4 files changed, 49 insertions(+), 12 deletions(-) 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 49af052b46ae8..418e7902b38bd 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 @@ -1476,21 +1476,21 @@ public void testGetTxnState() throws Exception { .build().get(); // test OPEN and TIMEOUT - assertEquals(transaction.getTxnState(), Transaction.State.OPEN); + assertEquals(transaction.getState(), Transaction.State.OPEN); Transaction timeoutTxn = transaction; - Awaitility.await().until(() -> timeoutTxn.getTxnState() == Transaction.State.TIMEOUT); + Awaitility.await().until(() -> timeoutTxn.getState() == Transaction.State.TIME_OUT); // test abort transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) .build().get(); transaction.abort().get(); - assertEquals(transaction.getTxnState(), Transaction.State.ABORTED); + assertEquals(transaction.getState(), Transaction.State.ABORTED); // test commit transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) .build().get(); transaction.commit().get(); - assertEquals(transaction.getTxnState(), Transaction.State.COMMITTED); + assertEquals(transaction.getState(), Transaction.State.COMMITTED); // test error transaction = pulsarClient.newTransaction().withTransactionTimeout(1, TimeUnit.SECONDS) @@ -1499,7 +1499,7 @@ public void testGetTxnState() throws Exception { .endTransaction(transaction.getTxnID(), 0, false); transaction.commit(); Transaction errorTxn = transaction; - Awaitility.await().until(() -> errorTxn.getTxnState() == Transaction.State.ERROR); + Awaitility.await().until(() -> errorTxn.getState() == Transaction.State.ERROR); // test committing transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) @@ -1507,7 +1507,7 @@ public void testGetTxnState() throws Exception { ((TransactionImpl) transaction).registerSendOp(new CompletableFuture<>()); transaction.commit(); Transaction committingTxn = transaction; - Awaitility.await().until(() -> committingTxn.getTxnState() == Transaction.State.COMMITTING); + Awaitility.await().until(() -> committingTxn.getState() == Transaction.State.COMMITTING); // test aborting transaction = pulsarClient.newTransaction().withTransactionTimeout(3, TimeUnit.SECONDS) @@ -1515,6 +1515,6 @@ public void testGetTxnState() throws Exception { ((TransactionImpl) transaction).registerSendOp(new CompletableFuture<>()); transaction.abort(); Transaction abortingTxn = transaction; - Awaitility.await().until(() -> abortingTxn.getTxnState() == Transaction.State.ABORTING); + Awaitility.await().until(() -> abortingTxn.getState() == Transaction.State.ABORTING); } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java index 37d9eb6967da4..3705607c7f923 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/TransactionEndToEndTest.java @@ -1050,7 +1050,7 @@ public void testTxnTimeOutInClient() throws Exception{ .build().get(); producer.newMessage().send(); Awaitility.await().untilAsserted(() -> { - Assert.assertEquals(((TransactionImpl)transaction).getState(), TransactionImpl.State.TIMEOUT); + Assert.assertEquals(((TransactionImpl)transaction).getState(), TransactionImpl.State.TIME_OUT); }); try { diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java index 71cab13181bff..4f45e407b6fb9 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java @@ -30,13 +30,50 @@ public interface Transaction { enum State { + + /** + * When the transaction is in the `OPEN` state, it can produce with transaction and ack with the transaction. + * + * When the transaction is in the `OPEN` state, it can commit or abort. + */ OPEN, + + /** + * When the client invokes commit, the state will change to `COMMITTING` from `OPEN`. + */ COMMITTING, + + /** + * When the client invokes abort, the state will change to `ABORTING` from `OPEN`. + */ ABORTING, + + /** + * When the client receives the response to the commit, the state will change to `COMMITTED` from `COMMITTING`. + */ COMMITTED, + + /** + * When the client receives the response to the abort, the state will change to `ABORTED` from `ABORTING`. + */ ABORTED, + + /** + * When the client invokes commit or abort but transaction not exist in coordinator, + * the state will change to `ERROR`. + * + * When the client invokes commit, but the transaction state in coordinator is committed or committing, + * the state will change to `ERROR`. + * + * When the client invokes abort, but the transaction state in coordinator is aborted or aborting, + * the state will change to `ERROR`. + */ ERROR, - TIMEOUT + + /** + * When the transaction timeout and the state is in `OPEN`, the state will change to `TIME_OUT` from `OPEN`. + */ + TIME_OUT } /** @@ -64,6 +101,6 @@ enum State { * * @return {@link State} the state of the transaction. */ - State getTxnState(); + State getState(); } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java index 0c0bcca07d623..833b0957d1c8a 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/transaction/TransactionImpl.java @@ -70,7 +70,7 @@ public class TransactionImpl implements Transaction , TimerTask { @Override public void run(Timeout timeout) throws Exception { - STATE_UPDATE.compareAndSet(this, State.OPEN, State.TIMEOUT); + STATE_UPDATE.compareAndSet(this, State.OPEN, State.TIME_OUT); } TransactionImpl(PulsarClientImpl client, @@ -206,7 +206,7 @@ public TxnID getTxnID() { } @Override - public State getTxnState() { + public State getState() { return state; } From 07cea27c82677debf5966f6d82daa401832551b0 Mon Sep 17 00:00:00 2001 From: congbo <39078850+congbobo184@users.noreply.github.com> Date: Mon, 5 Sep 2022 21:08:35 +0800 Subject: [PATCH 3/4] Update pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java Co-authored-by: Anonymitaet <50226895+Anonymitaet@users.noreply.github.com> --- .../org/apache/pulsar/client/api/transaction/Transaction.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java index 4f45e407b6fb9..77cba027a58b1 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java @@ -34,7 +34,7 @@ enum State { /** * When the transaction is in the `OPEN` state, it can produce with transaction and ack with the transaction. * - * When the transaction is in the `OPEN` state, it can commit or abort. + * When a transaction is in the `OPEN` state, it can commit or abort. */ OPEN, From d4594e7625e7de4d6a692dbf5299d39132353f69 Mon Sep 17 00:00:00 2001 From: congbobo184 Date: Mon, 5 Sep 2022 21:17:20 +0800 Subject: [PATCH 4/4] fix some doc comment --- .../client/api/transaction/Transaction.java | 26 ++++++++++--------- 1 file changed, 14 insertions(+), 12 deletions(-) diff --git a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java index 4f45e407b6fb9..453dc12cb42e1 100644 --- a/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java +++ b/pulsar-client-api/src/main/java/org/apache/pulsar/client/api/transaction/Transaction.java @@ -32,46 +32,48 @@ public interface Transaction { enum State { /** - * When the transaction is in the `OPEN` state, it can produce with transaction and ack with the transaction. + * When a transaction is in the `OPEN` state, messages can be produced and acked with this transaction. * * When the transaction is in the `OPEN` state, it can commit or abort. */ OPEN, /** - * When the client invokes commit, the state will change to `COMMITTING` from `OPEN`. + * When a client invokes a commit, the transaction state is changed from `OPEN` to `COMMITTING`. */ COMMITTING, /** - * When the client invokes abort, the state will change to `ABORTING` from `OPEN`. + * When a client invokes an abort, the transaction state is changed from `OPEN` to `ABORTING`. */ ABORTING, /** - * When the client receives the response to the commit, the state will change to `COMMITTED` from `COMMITTING`. + * When a client receives a response to a commit, the transaction state is changed from + * `COMMITTING` to `COMMITTED`. */ COMMITTED, /** - * When the client receives the response to the abort, the state will change to `ABORTED` from `ABORTING`. + * When a client receives a response to an abort, the transaction state is changed from `ABORTING` to `ABORTED`. */ ABORTED, /** - * When the client invokes commit or abort but transaction not exist in coordinator, - * the state will change to `ERROR`. + * When a client invokes a commit or an abort, but a transaction does not exist in a coordinator, + * then the state is changed to `ERROR`. * - * When the client invokes commit, but the transaction state in coordinator is committed or committing, - * the state will change to `ERROR`. + * When a client invokes a commit, but the transaction state in a coordinator is `ABORTED` or `ABORTING`, + * then the state is changed to `ERROR`. * - * When the client invokes abort, but the transaction state in coordinator is aborted or aborting, - * the state will change to `ERROR`. + * When a client invokes an abort, but the transaction state in a coordinator is `COMMITTED` or `COMMITTING`, + * then the state is changed to `ERROR`. */ ERROR, /** - * When the transaction timeout and the state is in `OPEN`, the state will change to `TIME_OUT` from `OPEN`. + * When a transaction is timed out and the transaction state is `OPEN`, + * then the transaction state is changed from `OPEN` to `TIME_OUT`. */ TIME_OUT }