diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java index 364ca181719c8..c0701c5cb389f 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/SenderTest.java @@ -32,6 +32,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.errors.ClusterAuthorizationException; import org.apache.kafka.common.errors.InvalidRequestException; +import org.apache.kafka.common.errors.InvalidTxnStateException; import org.apache.kafka.common.errors.NetworkException; import org.apache.kafka.common.errors.RecordTooLargeException; import org.apache.kafka.common.errors.TimeoutException; @@ -3085,7 +3086,7 @@ public void testReceiveFailedBatchTwiceWithTransactions() throws Exception { time.sleep(20); - sendIdempotentProducerResponse(0, tp0, Errors.INVALID_RECORD, -1); + sendIdempotentProducerResponse(0, tp0, Errors.INVALID_TXN_STATE, -1); sender.runOnce(); // receive late response // Loop once and confirm that the transaction manager does not enter a fatal error state @@ -3107,6 +3108,45 @@ public void testReceiveFailedBatchTwiceWithTransactions() throws Exception { txnManager.beginTransaction(); } + @Test + public void testInvalidTxnStateIsAnAbortableError() throws Exception { + ProducerIdAndEpoch producerIdAndEpoch = new ProducerIdAndEpoch(123456L, (short) 0); + apiVersions.update("0", NodeApiVersions.create(ApiKeys.INIT_PRODUCER_ID.id, (short) 0, (short) 3)); + TransactionManager txnManager = new TransactionManager(logContext, "testInvalidTxnState", 60000, 100, apiVersions); + + setupWithTransactionState(txnManager); + doInitTransactions(txnManager, producerIdAndEpoch); + + txnManager.beginTransaction(); + txnManager.maybeAddPartition(tp0); + client.prepareResponse(buildAddPartitionsToTxnResponseData(0, Collections.singletonMap(tp0, Errors.NONE))); + sender.runOnce(); + + Future request = appendToAccumulator(tp0); + sender.runOnce(); // send request + sendIdempotentProducerResponse(0, tp0, Errors.INVALID_TXN_STATE, -1); + + // Return InvalidTxnState error. It should be abortable. + sender.runOnce(); + assertFutureFailure(request, InvalidTxnStateException.class); + assertTrue(txnManager.hasAbortableError()); + TransactionalRequestResult result = txnManager.beginAbort(); + sender.runOnce(); + + // Once the transaction is aborted, we should be able to begin a new one. + respondToEndTxn(Errors.NONE); + sender.runOnce(); + assertTrue(txnManager::isInitializing); + prepareInitProducerResponse(Errors.NONE, producerIdAndEpoch.producerId, producerIdAndEpoch.epoch); + sender.runOnce(); + assertTrue(txnManager::isReady); + + assertTrue(result.isSuccessful()); + result.await(); + + txnManager.beginTransaction(); + } + private void verifyErrorMessage(ProduceResponse response, String expectedMessage) throws Exception { Future future = appendToAccumulator(tp0, 0L, "key", "value"); sender.runOnce(); // connect diff --git a/core/src/main/scala/kafka/log/UnifiedLog.scala b/core/src/main/scala/kafka/log/UnifiedLog.scala index 2ddaa284fcd12..028f2fbbd5e48 100644 --- a/core/src/main/scala/kafka/log/UnifiedLog.scala +++ b/core/src/main/scala/kafka/log/UnifiedLog.scala @@ -1031,7 +1031,7 @@ class UnifiedLog(@volatile var logStartOffset: Long, // ongoing. If the transaction is expected to be ongoing, we will not set a verification guard. If the transaction is aborted, hasOngoingTransaction is false and // requestVerificationGuard is null, so we will throw an error. A subsequent produce request (retry) should create verification state and return to phase 1. if (batch.isTransactional && !hasOngoingTransaction(batch.producerId) && batchMissingRequiredVerification(batch, requestVerificationGuard)) - throw new InvalidRecordException("Record was not part of an ongoing transaction") + throw new InvalidTxnStateException("Record was not part of an ongoing transaction") } // We cache offset metadata for the start of each transaction. This allows us to diff --git a/core/src/main/scala/kafka/server/AddPartitionsToTxnManager.scala b/core/src/main/scala/kafka/server/AddPartitionsToTxnManager.scala index fc5705042fe47..05e40014669e7 100644 --- a/core/src/main/scala/kafka/server/AddPartitionsToTxnManager.scala +++ b/core/src/main/scala/kafka/server/AddPartitionsToTxnManager.scala @@ -142,9 +142,9 @@ class AddPartitionsToTxnManager(config: KafkaConfig, client: NetworkClient, time if (addPartitionsToTxnResponseData.errorCode != 0) { error(s"AddPartitionsToTxnRequest for node ${response.destination} returned with error ${Errors.forCode(addPartitionsToTxnResponseData.errorCode)}.") // The client should not be exposed to CLUSTER_AUTHORIZATION_FAILED so modify the error to signify the verification did not complete. - // Older clients return with INVALID_RECORD and newer ones can return with INVALID_TXN_STATE. + // Return INVALID_TXN_STATE. val finalError = if (addPartitionsToTxnResponseData.errorCode == Errors.CLUSTER_AUTHORIZATION_FAILED.code) - Errors.INVALID_RECORD.code + Errors.INVALID_TXN_STATE.code else addPartitionsToTxnResponseData.errorCode @@ -160,9 +160,6 @@ class AddPartitionsToTxnManager(config: KafkaConfig, client: NetworkClient, time val code = if (partitionResult.partitionErrorCode == Errors.PRODUCER_FENCED.code) Errors.INVALID_PRODUCER_EPOCH.code - // Older clients return INVALID_RECORD - else if (partitionResult.partitionErrorCode == Errors.INVALID_TXN_STATE.code) - Errors.INVALID_RECORD.code else partitionResult.partitionErrorCode unverified.put(tp, Errors.forCode(code)) diff --git a/core/src/main/scala/kafka/server/ReplicaManager.scala b/core/src/main/scala/kafka/server/ReplicaManager.scala index 69d291f84fa7b..d56762a151d6c 100644 --- a/core/src/main/scala/kafka/server/ReplicaManager.scala +++ b/core/src/main/scala/kafka/server/ReplicaManager.scala @@ -732,8 +732,7 @@ class ReplicaManager(val config: KafkaConfig, debug("Produce to local log in %d ms".format(time.milliseconds - sTime)) val unverifiedResults = unverifiedEntries.map { case (topicPartition, error) => - // NOTE: Older clients return INVALID_RECORD, but newer clients will return INVALID_TXN_STATE - val message = if (error.equals(Errors.INVALID_RECORD)) "Partition was not added to the transaction" else error.message() + val message = if (error == Errors.INVALID_TXN_STATE) "Partition was not added to the transaction" else error.message() topicPartition -> LogAppendResult( LogAppendInfo.UNKNOWN_LOG_APPEND_INFO, Some(error.exception(message)) diff --git a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala index b451dd8578416..08bc2e7785ee8 100644 --- a/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala +++ b/core/src/test/scala/unit/kafka/cluster/PartitionTest.scala @@ -24,7 +24,7 @@ import kafka.server._ import kafka.server.checkpoints.OffsetCheckpoints import kafka.utils._ import kafka.zk.KafkaZkClient -import org.apache.kafka.common.errors.{ApiException, FencedLeaderEpochException, InconsistentTopicIdException, NotLeaderOrFollowerException, OffsetNotAvailableException, OffsetOutOfRangeException, UnknownLeaderEpochException} +import org.apache.kafka.common.errors.{ApiException, FencedLeaderEpochException, InconsistentTopicIdException, InvalidTxnStateException, NotLeaderOrFollowerException, OffsetNotAvailableException, OffsetOutOfRangeException, UnknownLeaderEpochException} import org.apache.kafka.common.message.{AlterPartitionResponseData, FetchResponseData} import org.apache.kafka.common.message.LeaderAndIsrRequestData.LeaderAndIsrPartitionState import org.apache.kafka.common.protocol.{ApiKeys, Errors} @@ -32,7 +32,7 @@ import org.apache.kafka.common.record.FileRecords.TimestampAndOffset import org.apache.kafka.common.record._ import org.apache.kafka.common.requests.{AlterPartitionResponse, FetchRequest, ListOffsetsRequest, RequestHeader} import org.apache.kafka.common.utils.SystemTime -import org.apache.kafka.common.{InvalidRecordException, IsolationLevel, TopicPartition, Uuid} +import org.apache.kafka.common.{IsolationLevel, TopicPartition, Uuid} import org.apache.kafka.metadata.LeaderRecoveryState import org.junit.jupiter.api.Assertions._ import org.junit.jupiter.api.Test @@ -3387,14 +3387,14 @@ class PartitionTest extends AbstractPartitionTest { producerId = producerId) // When verification guard is not there, we should not be able to append. - assertThrows(classOf[InvalidRecordException], () => partition.appendRecordsToLeader(transactionRecords(), origin = AppendOrigin.CLIENT, requiredAcks = 1, RequestLocal.withThreadConfinedCaching)) + assertThrows(classOf[InvalidTxnStateException], () => partition.appendRecordsToLeader(transactionRecords(), origin = AppendOrigin.CLIENT, requiredAcks = 1, RequestLocal.withThreadConfinedCaching)) // Before appendRecordsToLeader is called, ReplicaManager will call maybeStartTransactionVerification. We should get a non-null verification object. val verificationGuard = partition.maybeStartTransactionVerification(producerId, 3, 0) assertNotNull(verificationGuard) // With the wrong verification guard, append should fail. - assertThrows(classOf[InvalidRecordException], () => partition.appendRecordsToLeader(transactionRecords(), + assertThrows(classOf[InvalidTxnStateException], () => partition.appendRecordsToLeader(transactionRecords(), origin = AppendOrigin.CLIENT, requiredAcks = 1, RequestLocal.withThreadConfinedCaching, Optional.of(new Object))) // We should return the same verification object when we still need to verify. Append should proceed. diff --git a/core/src/test/scala/unit/kafka/server/AddPartitionsToTxnManagerTest.scala b/core/src/test/scala/unit/kafka/server/AddPartitionsToTxnManagerTest.scala index 232e9d012d905..9231fdc124f4b 100644 --- a/core/src/test/scala/unit/kafka/server/AddPartitionsToTxnManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/AddPartitionsToTxnManagerTest.scala @@ -201,7 +201,7 @@ class AddPartitionsToTxnManagerTest { assertEquals(expectedDisconnectedErrors, transaction1Errors) assertEquals(expectedDisconnectedErrors, transaction2Errors) - val expectedTopLevelErrors = topicPartitions.map(_ -> Errors.INVALID_RECORD).toMap + val expectedTopLevelErrors = topicPartitions.map(_ -> Errors.INVALID_TXN_STATE).toMap val topLevelErrorAddPartitionsResponse = new AddPartitionsToTxnResponse(new AddPartitionsToTxnResponseData().setErrorCode(Errors.CLUSTER_AUTHORIZATION_FAILED.code())) val topLevelErrorResponse = clientResponse(topLevelErrorAddPartitionsResponse) addTransactionsToVerify() @@ -212,9 +212,9 @@ class AddPartitionsToTxnManagerTest { val preConvertedTransaction1Errors = topicPartitions.map(_ -> Errors.PRODUCER_FENCED).toMap val expectedTransaction1Errors = topicPartitions.map(_ -> Errors.INVALID_PRODUCER_EPOCH).toMap val preConvertedTransaction2Errors = Map(new TopicPartition("foo", 1) -> Errors.NONE, - new TopicPartition("foo", 2) -> Errors.INVALID_RECORD, + new TopicPartition("foo", 2) -> Errors.INVALID_TXN_STATE, new TopicPartition("foo", 3) -> Errors.NONE) - val expectedTransaction2Errors = Map(new TopicPartition("foo", 2) -> Errors.INVALID_RECORD) + val expectedTransaction2Errors = Map(new TopicPartition("foo", 2) -> Errors.INVALID_TXN_STATE) val transaction1ErrorResponse = AddPartitionsToTxnResponse.resultForTransaction(transactionalId1, preConvertedTransaction1Errors.asJava) val transaction2ErrorResponse = AddPartitionsToTxnResponse.resultForTransaction(transactionalId2, preConvertedTransaction2Errors.asJava) diff --git a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala index 95de4b83709c7..7248330741100 100644 --- a/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala +++ b/core/src/test/scala/unit/kafka/server/ReplicaManagerTest.scala @@ -2231,8 +2231,8 @@ class ReplicaManagerTest { // Confirm we did not write to the log and instead returned error. val callback: AddPartitionsToTxnManager.AppendCallback = appendCallback.getValue() - callback(Map(tp0 -> Errors.INVALID_RECORD).toMap) - assertEquals(Errors.INVALID_RECORD, result.assertFired.error) + callback(Map(tp0 -> Errors.INVALID_TXN_STATE).toMap) + assertEquals(Errors.INVALID_TXN_STATE, result.assertFired.error) assertEquals(verificationGuard, getVerificationGuard(replicaManager, tp0, producerId)) // This time verification is successful.